락으로 막지 말고, 경합이 생기지 않게 나눠라

작성 · 수정


출처: 카카오페이 기술블로그 — 지연이체 서비스 개발기: 은행 점검 시간 끝나면 송금해 드릴게요! (feat. 발표 후기)

if(kakaoAI)2024에서 카카오페이 엘모가 발표한 내용을 정리한 글이다. 사내 MQ 통합으로 RabbitMQ가 사라지게 되자 지연이체를 Kafka 위에서 다시 만든 기록이다. 원문을 읽고 정리한 내용과 내 생각을 함께 적는다.

문제

지연이체는 은행 점검 시간에 송금을 시도한 사용자가 예약을 걸어두면 점검이 끝난 뒤 자동으로 송금해 주는 기능이다. 기존 구조는 RabbitMQ의 지연 큐였다. 등록 시점에 은행 펌뱅킹으로 점검 완료 시간을 받아 x분 뒤로 계산해 지연 큐에 넣으면, 시간이 지나 실행 큐로 옮겨지고 Consumer가 송금을 실행한다. 지연이라는 요구사항이 브로커 기능 하나로 그대로 충족되는 구조다.

이 구조를 걷어낸 건 기술적 결함 때문이 아니었다. 사내에서 공용 MQ를 Kafka 하나로 통합하기로 결정하면서 RabbitMQ 운영이 중단될 예정이었다. 서비스를 멈출 수는 없으니 재설계 외에 선택지가 없었다.

여기서부터 제약의 성격이 달라진다. 보통 아키텍처 변경은 지금 구조가 감당을 못 해서 시작하는데, 이 경우는 잘 돌던 구조를 조직 표준 때문에 들어내는 일이다. 그래서 질문이 “더 좋은가”에서 “원래 기능을 브로커 도움 없이 어떻게 복원하는가”로 바뀐다. Kafka에는 RabbitMQ 같은 지연 큐가 없다.

분석

원문은 후보 두 가지를 놓고 하나씩 탈락시킨다.

스케줄러 단독. 5분마다 상태가 DELAY이고 실행 시각이 된 건을 읽어 송금을 요청한다. 가장 단순한 답이고 실제로 1단계 구현이 이랬다. 문제는 한 번에 수천 건 이상이 걸린다는 것이었다. 점검 종료라는 특정 시각에 등록 건이 몰려 있으니 평균이 아니라 스파이크가 기준이 된다. 5분 안에 수천 건을 끝내는 건 그 구조로는 불가능했다.

스케줄러 여러 대. 분산시키면 되지 않느냐는 자연스러운 다음 수인데, 원문은 세 가지 이유로 접는다. 서로 겹치지 않게 읽으려면 복잡한 분기가 붙고, 인프라 관리가 늘고, 여러 대가 동시에 송금 실행 요청을 쏘면 송금 실행 서버가 부담을 받는다.

마지막 이유가 중요하다. 스케줄러를 늘리는 건 내 쪽 처리량만 올리는 일이고, 그 부하는 고스란히 하류로 간다. 분산은 했는데 조절 지점이 없는 구조다.

그래서 남은 답이 Kafka다. 토픽 transfer-delay 하나에 파티션 3개, Consumer 2대로 시작했다. 스케줄러는 5분마다 대상 건을 읽어 Produce만 하고, 실제 송금 실행은 Consumer가 맡는다. 스케줄러 여러 대 방안과 달리 조회는 스케줄러, 실행은 Consumer로 나뉘면서 실행 쪽 병렬도를 Consumer 수와 설정으로 따로 조절할 수 있게 됐다.

해결

이 글에서 가장 볼 만한 부분은 성능 튜닝이 아니라 중복 방지다. 한 번에 해결된 게 아니라, 막으면 다음 문제가 나오는 연쇄였다.

1. 같은 송금 건이 토픽에 중복으로 쌓인다. 스케줄러는 5분마다 도는데 이전 회차에 Produce한 건이 아직 소비되지 않았으면 다음 회차가 같은 건을 또 읽어 간다. 처음 붙인 대응은 Consumer 쪽 상태 체크였다. DELAY가 아니면 건너뛴다. 최종적으로는 스케줄러가 Produce 직전에 상태를 PREPARATION으로 바꾸고 보내도록 바꿨다. 앞의 것은 잘못 들어온 메시지를 소비 시점에 걸러내는 방식이고, 뒤의 것은 애초에 두 번 들어가지 않게 하는 방식이다. 방어선을 뒤에서 앞으로 옮겼다.

2. Consumer가 동시에 같은 건을 소비한다. 상태 플래그는 스케줄러 회차 사이의 중복을 막을 뿐, 이미 토픽에 들어간 뒤 여러 Consumer가 동시에 집는 상황은 남는다. 여기에는 유저락을 걸었다. 동시에 같은 건을 실행하려 해도 한쪽만 통과한다.

3. 락이 만든 연쇄 실패. 한 사용자가 여러 건을 예약했는데 그 건들이 서로 다른 Consumer로 흩어지면, 유저락 때문에 한 건만 성공하고 나머지가 실패한다. 중복을 막으려고 넣은 장치가 정상 요청까지 떨어뜨린 것이다.

마지막 해법은 앞의 둘과 성격이 다르다. 락을 더 정교하게 다듬지 않았다. Produce할 때 Record Key에 userId를 넣었다. 같은 사용자의 건은 같은 파티션으로 가고, 한 파티션은 한 Consumer가 가져가니, 그 건들은 한 Consumer에서 순차 처리된다.

락은 경합이 일어난 뒤에 조정하는 장치이고, 파티션 키는 경합 자체가 생기지 않게 배치를 바꾸는 장치다. 앞의 두 단계가 “충돌하면 어떻게 할까”였다면 마지막은 “충돌할 수 있는 것들을 애초에 한 줄에 세우자”다. 유저락은 그대로 남아 있지만 이제 그건 정상 경로가 아니라 최후의 안전망이다.

속도 개선은 세 가지를 같이 올렸다. Consumer를 2대에서 3대로 늘리고, max.poll.records를 1에서 20으로 올려 배치로 읽고, ForkJoinPool 스레드를 1에서 10으로 늘려 배치 안을 병렬로 돌렸다.

val groupedMessages = delayMessages.groupBy { it.userId }

forkJoinPool.submit {
    groupedMessages.values.parallelStream().forEach { messages ->
        delaySend(headers, messages)
    }
}.join()

멀티스레드를 넣으면서도 userId로 묶어 그룹 단위로만 병렬화한다는 점을 봐야 한다. 파티션 키로 만들어 둔 “같은 유저는 순차” 보장을 스레드가 다시 깨뜨리면 3번 문제가 Consumer 안에서 그대로 재현된다. 파티션 수준에서 세운 규칙을 스레드 수준에서도 똑같이 유지한 셈이다.

파티션은 3개 그대로 뒀다. 원문이 든 이유는 한 번 늘린 파티션 수는 줄이려 하면 InvalidPartitionsException이 나기 때문이다. 여기에 하나 더 붙일 이유가 있다. 순서 보장을 파티션 키에 의존하기 시작한 순간부터 파티션 수 변경은 순서 보장을 깨는 변경이 된다. 키를 파티션에 매핑하는 계산에 파티션 개수가 들어가므로, 개수가 바뀌면 같은 userId가 다른 파티션으로 향한다. 기존 파티션에 남아 있던 건과 새로 들어가는 건이 서로 다른 Consumer에서 동시에 처리될 수 있다. Consumer만 늘리고 파티션을 유지한 선택은 이 맥락에서 합리적이다. 다만 Consumer 3대와 파티션 3개는 이미 병렬성 상한이다. 이 방향으로 더 가려면 파티션을 건드려야 하고, 그건 방금 말한 비용을 치르는 일이다.

결과

항목개선 전개선 후
Consumer 대수23
Read Records 수120
Thread 수110
1분당 처리 건수91건728건
실행 시간68분8분

원문 표현으로 1분당 처리량 800% 증가, 실행 속도 8배다.

곱해보면 이 숫자가 더 흥미로워진다. Consumer 1.5배에 배치 20배, 스레드 10배를 곱하면 산술적으로는 300배다. 실제로는 8배다. 그 차이가 어디로 갔는지 생각해 볼 만하다.

답은 병목이 Consumer 쪽에 있지 않았다는 것이다. 송금 실행 서버, 그 뒤의 은행 API, DB 중 어딘가가 상한을 잡고 있고 Consumer는 그 상한까지 밀어 올린 것에 가깝다. 곱셈이 그대로 나오지 않는 게 오히려 정상이다. 그대로 나왔다면 아직 하류에 여유가 남아 있다는 뜻이고, 계속 올리다 보면 어느 지점에서 실행 서버가 먼저 무너진다.

원문이 max.poll.records와 스레드 수를 “송금 실행 내부 서버 부담과 목표 송금 실행 시간을 고려하여” 정했다고 적은 게 그래서 눈에 남는다. 이 값들은 최대치를 찾는 실험의 결과가 아니라 하류가 견디는 한계와 맞춰야 할 목표 시간 사이에서 고른 값이다. 처리량 튜닝은 내 쪽 숫자를 올리는 일이 아니라, 하류가 감당하는 선까지만 올리고 멈추는 일이다. 스케줄러 여러 대 방안을 접은 이유도 결국 하류 부담이었다.

이 블로그에는 Kafka를 걷어내고 DB 작업 큐로 글이 있다. 브로커를 지우고 스케줄러와 DB 테이블로 돌아간 기록이니 방향이 정반대다. 둘 중 하나가 틀린 게 아니라 조건이 다르다.

카카오페이 쪽에는 브로커 값어치가 나오는 조건이 두 개 있었다. 사내 MQ를 Kafka로 통합한다는 결정이 이미 외부에서 내려와 있었고, 은행 점검 종료라는 특정 시각에 수천 건 이상이 한꺼번에 몰리는 스파이크가 실재했다. 앞의 것은 기술 선택 이전의 제약이고, 뒤의 것은 유입이 처리 능력을 앞지르는 물량을 어딘가에 안전하게 눕혀 두고 Consumer 수로 나눠 가져야 하는 상황이다. 여기에 파티션 키로 순서를 보장해야 할 요구까지 있었다. Kafka가 하는 일이 명확하다.

내 쪽 사례에는 그 조건이 하나도 없었다. 트래픽은 스파이크를 만들지 않았고, 처리량 상한은 브로커가 아니라 건당 몇 초짜리 외부 호출이 정했고, 조직 표준이 브로커를 강제하지도 않았다. 같은 도구가 한쪽에서는 병목을 흡수하고 다른 쪽에서는 정합성 방어 코드만 늘렸다.

인상 깊었던 점

중복 방지 세 단계를 한 흐름으로 보여준 게 이 글의 가치다. 결과만 요약하면 “상태 플래그와 파티션 키로 중복을 막았다” 한 줄인데, 그렇게 적으면 파티션 키가 왜 나왔는지가 사라진다. 락을 걸어봐야 락 때문에 생기는 실패를 만나고, 그제야 문제를 경합 해소가 아니라 경합 회피로 다시 정의하게 된다. 중간 단계는 실패해서 버린 게 아니라, 지나야 다음 문제가 보이는 단계였다.

동시성 문제를 만나면 나는 대체로 락부터 떠올린다. 락은 이미 충돌한 둘 중 하나를 막아 세울 뿐이고, 막힌 쪽을 어떻게 할지는 그대로 남는다. 재시도할지, 실패로 볼지, 대기시킬지. 지연이체에서 3번 문제가 정확히 그 자리에서 나왔다. 애초에 같은 자원에 손대는 것들을 한 실행 단위로 몰아둘 수 있는지 먼저 보는 게 순서다. 파티션 키로 묶든, 같은 스레드로 묶든, 같은 트랜잭션으로 묶든.

그리고 이 글의 결론 역시 카카오페이의 조건 위에서 나왔다. RabbitMQ를 걷어낸 건 성능 문제가 아니라 조직 결정이었고, Kafka를 고른 건 그 결정 안에서 스파이크를 감당할 수 있는 답이 그것이었기 때문이다. 조건이 바뀌면 답도 바뀐다.