콘서트 티켓팅 혹은 성수기의 항공기, 호텔 예약 등이 시작될 때, 많은 사용자들이 갑자기 동시 접속을 하게 된다. 이 경우, 중요한 것은 바로 아래 두 가지다.
첫째, 서버가 트래픽을 감당할 수 있는지 여부다.
둘째, 유입된 사용자의 순서를 지켜서 사용자들을 줄세우고, 순차적으로 오픈해줄 수 있는지다.
이 글에서는 첫 번째 내용만 자세히 설명하려고 한다.
다음 글에서 두 번째 내용을 어떻게 해결했는지 작성할 것이다.
트래픽이 단일 애플리케이션 수준에서 처리할 수 없게되면, 서버를 수평 확장(스케일아웃)함으로써 문제를 해결할 수 있다. 하지만 비용이 들 뿐더러, 갑작스러운 트래픽 유입의 경우, 스케일아웃으로 빠른 대응은 가능하지만 즉각적인 대응은 불가능하다. 왜냐하면 서버 실행에도 시간이 걸리기 때문이다.
Spring Webflux (출처)
기존 스프링 프레임워크에 내장되어 있던 웹 프레임워크인 Spring Web MVC는 서블릿 API와 컨테이너를 목적으로 설계되었다. 리액티브 스택의 웹 프레임워크인 Spring Webflux는 스프링 5버전에서 추가된 것으로, 완전히 논블락킹이고, 리액티브 스트림즈의 백프레셔를 지원하고, Netty, Undertow, Servlet 컨테이너 등의 서버에서 실행된다.
기존 MVC 방식에서는 OS의 스레드 기반으로 웹 요청을 처리하고, 100개의 요청이 들어왔을 때 100개의 스레드를 모두 사용하는 형태로 동작한다. 그렇기 때문에, 대용량 트래픽 처리에서 메모리 사용량이 급증하는 형태를 보인다.
또한 외부 DB와의 IO 발생 시, 블락킹까지 일어난다면 처리 성능에 더욱 부하가 생긴다.
그래서 스케일 아웃 전략을 취하기 전, 애플리케이션 서버 자체의 처리량을 높이기 위한 시도로 Spring Webflux를 사용해보기로 했다.
Spring Webflux는 네티 기반으로 네트워크 비동기 IO 작업을 처리한 후, 이벤트 루프 기반으로 작업을 빠르게 처리한다. 또한, CPU 코어 갯수가 4개 이상이어도 일반적으로 4개의 스레드를 만들고 이벤트 루프 방식으로 요청을 처리한다.
Redis는 인메모리 데이터 저장소이기 때문에 처리 속도가 빠르다. 단일 스레드로 처리하며, 버전 6부터는 외부 요청은 따로 이벤트 루프 방식으로 처리하기 때문에, 속도도 더 빨라졌다. persistent on disk 설정으로 데이터의 영구 저장도 할 수 있어 메모리 기반 저장소라는 단점도 보완되었다.
하지만 애플리케이션 서버 입장에서 Redis에 요청을 보내는 것은 네트워크 IO 작업이기 때문에, 분명 블락킹이 일어날 수 있는 지점이기도 하다. 그래서 Spring Data Redis에서는 ReactiveRedisTemplate을 지원하며, 이를 통해 WebFlux 프레임워크를 사용하는 애플리케이션 서버는 API 요청을 완전히 논블락킹으로 처리할 수 있게 된다.
결론
Spring WebFlux로 사용자 대기열 처리 서버를 구성하고,ReactiveRedisTemplate을 사용해서 Redis와의 통신도 논블락킹으로 처리하려고 한다.
public Mono<AddToQueueInfo> addToQueue() {
var unixTimestamp = System.currentTimeMillis();
var uuid = UUID.randomUUID();
log.info("ADDING TO QUEUE... timestamp = {}, uuid = {}", unixTimestamp, uuid);
return reactiveRedisTemplate.opsForZSet().add(USER_QUEUE_WAIT_KEY, String.valueOf(uuid), unixTimestamp)
.doOnNext(user -> log.info("UUID : {}, THREAD : {}", uuid, Thread.currentThread().getName()))
.filter(user -> user)
.flatMap(i -> reactiveRedisTemplate.opsForZSet().rank(USER_QUEUE_WAIT_KEY, uuid.toString())
.map(rank -> {
log.info("USER RANK IN THE WAITING QUEUE, RANK : {}, UUID : {}, THREAD : {}", rank + 1, uuid, Thread.currentThread().getName());
return new AddToQueueInfo(rank + 1, uuid.toString());
})
)
.onErrorResume(e -> Mono.error(FAILED_TO_ADD_IN_WAIT_QUEUE.build()));
}
@Scheduled(initialDelay = 5000, fixedDelay = 10000)
public void scheduleAllowUser() {
if (!scheduling) {
log.info("Scheduling is disabled...");
return;
}
log.info("Scheduling is called...");
reactiveRedisTemplate.scan(ScanOptions.scanOptions()
.match(USER_QUEUE_PROCEED_KEY_FOR_SCAN)
.build())
.count()
.doOnNext(size -> log.info("USER QUEUE IN PROGRESS SIZE : {}", size))
.flatMap(size -> {
long diff = USER_QUEUE_PROCEED_SIZE - size;
if (diff > 0) {
return reactiveRedisTemplate.opsForZSet().popMin(USER_QUEUE_WAIT_KEY, diff)
.flatMap(user -> reactiveRedisTemplate.opsForSet()
.add(USER_QUEUE_PROCEED_KEY.formatted(user.getValue()), user.getValue())
.then(reactiveRedisTemplate.expire(USER_QUEUE_PROCEED_KEY.formatted(user.getValue()), Duration.ofMinutes(10))))
.then();
} else {
return Mono.empty();
}
})
.subscribe();
}