카프카는 순서를 보장하는 게 아니라 보존한다

작성 · 수정


출처: 우아한형제들 기술블로그 — 우리 팀은 카프카를 어떻게 사용하고 있을까

하루 100만 건 이상의 배민배달을 중계하는 딜리버리서비스팀이 카프카를 세 자리에 쓰고 있는 방식을 소개한 글이다. 이벤트 브로커, 이벤트 버스, 그리고 실시간 분석 파이프라인. 원문을 읽고 정리한 내용과 내 생각을 함께 적는다.

문제

딜리버리서비스팀은 배민배달을 받아 여러 배달서비스 중 하나로 분배하고 그 과정을 중계한다. 주문이벤트를 받아 배달 프로세스를 관리하는 주문/배달서버와, 발행된 이벤트를 분석하는 분석서버로 나뉜 이벤트 기반 분산 시스템이다. 각 서버군은 여러 대로 구성된다.

이 구조에서 팀이 풀려는 문제는 세 가지다. 첫째, 배달 상태를 바꾸는 이벤트들이 발생한 순서대로 소비되어야 하고 그 과정에서 누락이 없어야 한다. 배차완료와 픽업준비요청처럼 거의 동시에 일어나는 이벤트가 뒤바뀌어 도착하면 컨슈머는 무엇이 먼저인지 알 수 없다. 둘째, 배달서버들이 인메모리로 들고 있는 분배 규칙을 운영자가 바꾸면 서버군 전체가 그 사실을 알아야 한다. 셋째, 배달 현황을 실시간 혹은 준실시간으로 보고 싶다.

세 문제에 모두 카프카를 썼다. 그런데 세 자리가 요구하는 정합성의 강도가 전부 다르다. 원문은 이 차이를 이름 붙여 말하지 않지만 각 경로가 어떤 약속을 하고 어떤 약속을 하지 않는지를 따라가 보면 이 글의 구조가 훨씬 선명해진다.

분석

순서를 정하는 것은 MySQL 커밋이고 카프카는 그것을 옮긴다

원문은 “카프카는 메시지 발행 순서에 따라 소비할 수 있도록 순서를 보장합니다”라고 쓴다. 문장 자체는 틀리지 않다. 다만 이 문장은 ‘발행 순서’를 이미 주어진 것으로 놓고 시작한다. 실제로 어려운 부분은 그 발행 순서가 어디서 정해지느냐다.

배차완료와 픽업준비요청을 다시 보자. 둘은 서로 다른 트랜잭션에서, 어쩌면 서로 다른 서버에서 거의 동시에 일어난다. 두 사건의 선후를 결정하는 것은 카프카가 아니다. 데이터베이스가 두 outbox insert를 커밋한 순서, 더 정확히는 그 두 커밋이 binlog에 기록된 위치다. 카프카는 순서를 만들어내지 않는다. 다른 곳에서 이미 정해진 순서를 한 번도 흐트러뜨리지 않고 옮길 뿐이다.

그 사슬은 이렇게 이어진다. 같은 키는 같은 outbox 테이블에 저장된다. 한 테이블에는 커넥터가 하나 붙는다. Debezium MySQL source connector는 태스크를 하나만 쓰도록 강제되어 있으니 그 커넥터 안에서 순서가 흐트러질 여지가 없다. 커넥터가 내보낸 레코드는 주문식별자나 배달식별자를 키로 갖고 같은 키는 같은 파티션으로 간다. 한 파티션은 한 컨슈머가 소비한다. 원문의 문장으로는 “같은 키는 같은 테이블에 저장되며, 한 테이블에서는 하나의 커넥터를 사용하기 때문에 같은 키에 대해서는 순서를 보장됩니다”다.

고리가 다섯 개인데 하나라도 빠지면 전체가 무너진다. 같은 키가 두 outbox 테이블에 흩어지면 커넥터 둘이 각자의 속도로 읽어 순서가 깨진다. 커넥터가 태스크를 둘 쓸 수 있었다면 같은 테이블 안에서도 깨진다. 키를 지정하지 않고 발행하면 파티션이 갈려 깨진다. 이 다섯 고리가 전부 한 방향을 가리키고 있다는 점이 이 설계의 핵심이고, 그건 우연으로 만들어지지 않는다.

그리고 이건 칭찬이다. 순서를 정하는 지점이 하나로 모여 있다는 뜻이니까. 분산 시스템에서 “무엇이 먼저인가”를 결정하는 자리가 여러 곳이면 그 자리들 사이의 합의가 다시 문제가 된다. 이 팀은 그 자리를 MySQL 커밋 하나로 밀어 넣었다.

”발행에 실패하면 롤백된다”는 문장은 인과가 뒤집혀 있다

원문에는 이런 문장이 있다. “메시지 발행에 실패하면 아웃박스테이블의 데이터도 롤백되기 때문에 하나의 트랜잭션으로 데이터 정합성을 관리하고 있습니다.”

CDC 기반 outbox에서 실제로 일어나는 일은 반대다. outbox insert는 비즈니스 데이터 변경과 같은 트랜잭션에서 커밋된다. 카프카 발행은 그보다 한참 뒤에 커넥터가 binlog를 읽어서 한다. 두 사건은 시간적으로도 분리돼 있고 트랜잭션도 공유하지 않는다. 커넥터가 발행에 실패했다고 해서 이미 커밋된 DB 행이 사라질 방법은 없다. 커넥터는 자기가 마지막으로 커밋한 오프셋부터 다시 읽어 재시도할 뿐이다.

이 패턴이 실제로 주는 약속은 “커밋됐으면 언젠가 반드시 발행된다”이지 “발행이 실패하면 커밋이 취소된다”가 아니다. 반대 방향은 성립한다. 비즈니스 트랜잭션이 롤백되면 같은 트랜잭션에 있던 outbox 행도 함께 사라지므로 커밋되지 않은 변경에 대한 이벤트는 나가지 않는다. 저자가 말하려던 것이 이쪽이었으리라 짐작한다. 어디까지나 추측이다.

문장을 바로잡고 나면 따라오는 게 있다. 재시도는 곧 중복이다. 커넥터가 재시작하면 마지막 커밋 오프셋 이후를 다시 내보내니 컨슈머는 같은 이벤트를 두 번 볼 수 있다. 원문도 “재시도 과정에서도 메시지의 순서는 보장되기를 바랐습니다”라고 적었는데, 그 바람이 이뤄지는 대가가 정확히 이것이다. 순서는 지켜지고 전달은 at-least-once가 된다.

그런데 원문 어디에도 컨슈머 멱등성 이야기가 없다. 배달상태를 바꾸는 이벤트를 두 번 받았을 때 무슨 일이 일어나는지, 이벤트 식별자로 중복을 거르는지, 상태 전이를 멱등하게 짰는지가 비어 있다. 순서 보장과 중복 허용은 한 몸이라서 앞을 얻은 시스템은 뒤를 반드시 다뤄야 한다. 실제로는 다루고 있을 가능성이 높다고 본다. 다만 글에는 없다.

커넥터를 늘려 얻은 처리량의 값은 DB가 치른다

단일 태스크라는 제약은 처리량 상한을 만든다. 팀은 토픽별로 outbox 테이블을 나누고 그걸 다시 식별자 기준으로 delivery-outbox1, delivery-outbox2, delivery-outbox3처럼 N개로 쪼갠 뒤 테이블마다 커넥터를 하나씩 붙였다. 합리적인 수평 분할이다.

다만 커넥터가 N개면 MySQL에 붙는 복제 클라이언트도 N개다. Debezium MySQL 커넥터는 특정 테이블만 골라 받는 것이 아니라 binlog 스트림을 받아서 자기 대상 테이블만 걸러낸다. 필터링이 커넥터 쪽에서 일어난다는 뜻이다. 그러면 커넥터를 하나 늘릴 때마다 DB는 binlog 전체를 한 번 더 내보낸다. 병렬화가 총 작업량을 줄이는 일은 거의 없다. 대개는 작업량을 다른 곳으로 옮긴다. 여기서는 커넥터의 처리 지연을 줄이는 대신 DB의 복제 출력과 네트워크를 N배로 늘렸다.

N을 어떻게 정했는지는 원문에 없다. 더 아쉬운 건 재분배 이야기가 없다는 점이다. 테이블을 하나 더 늘리면 키에서 테이블로 가는 매핑이 바뀐다. 매핑이 바뀌는 순간, 옮겨가는 키의 이벤트 일부는 옛 테이블에 남아 있고 나머지는 새 테이블로 들어간다. 두 테이블은 서로 다른 커넥터가 각자의 속도로 읽는다. 이 구간에서 순서 보장은 성립하지 않는다. 앞 절에서 본 다섯 고리 중 첫 번째가 잠시 끊기는 것이다. 무중단으로 넘기려면 옛 테이블이 완전히 소진될 때까지 새 테이블 쓰기를 미루는 식의 절차가 필요한데, 그 절차는 글에 없다.

정리 정책도 빈칸이다. outbox는 insert만 하는 테이블이라 계속 자란다. 하루 100만 건이 넘는 배달에 배달당 여러 이벤트가 붙으니 증가 속도가 작지 않다. CDC 방식에서는 binlog에 insert가 남기만 하면 되므로 행을 오래 보관할 이유가 없다. 같은 트랜잭션에서 insert 직후 delete 해버리는 방법이 알려져 있고 그렇게 하면 테이블은 사실상 비어 있는 채로 binlog만 흘러간다. 원문이 그렇게 하는지 아닌지는 알 수 없다.

이벤트에 값을 싣지 않은 것이 이벤트 버스에서 가장 잘한 판단이다

분배 규칙 변경 알림에는 Spring Cloud Bus의 RemoteApplicationEvent를 카프카 위에서 쓴다. DeliveryServiceRemoteApplicationEvent를 추상 클래스로 두고 그걸 상속한 RouteRuleRemoteEvent가 destination으로 “delivery”만 지정한다.

주목할 것은 RouteRuleRemoteEvent에 필드가 하나도 없다는 점이다. 바뀐 규칙을 실어 보내지 않는다. 받는 쪽 @EventListener가 하는 일은 routeRuleSetStore.load() 호출뿐이다. 상태를 전송하는 대신 “다시 읽어라”만 보내면 이벤트 순서와 유실과 지각 합류 문제가 전부 저장소 쪽으로 넘어간다. 규칙 변경 이벤트 두 개가 뒤바뀌어 도착해도 결과는 같다. 이벤트를 놓친 서버가 다음 이벤트에 반응하면 그 시점의 최신 규칙을 읽는다. 새로 뜬 서버는 anonymous 컨슈머 그룹이라 과거 이벤트를 보지 못하지만 기동하면서 저장소를 읽으니 애초에 볼 필요가 없다.

이벤트 버스 토픽이 파티션 하나인 것도 이 선택과 맞물린다. 순서에 의존하지 않는 알림이므로 처리량을 위해 파티션을 늘릴 이유가 없고 서버마다 컨슈머 그룹이 달라야 하니 파티션을 늘려봐야 각 그룹 안에서 놀고 있는 컨슈머만 생긴다.

사소한 지적이 하나 있다. 원문의 발행 측 코드는 routeRuleSetStore.load()를 부른 다음 remoteApplicationEventPublisher.publishEvent(new RouteRuleRemoteEvent())를 호출한다. Spring의 이벤트 발행은 로컬 리스너에도 전달되는 것이 보통이라 발행한 서버 자신이 handle()을 타면서 저장소를 한 번 더 읽을 가능성이 있다. Spring Cloud Bus가 자기 origin에서 온 이벤트를 어떻게 다루는지에 따라 달라지므로 단정하지는 않겠다. 어느 쪽이든 load()가 두 번 도는 것뿐이라 해가 없다.

anonymous 그룹을 모니터링에서 걸러낸다는 대목은 운영 냄새가 조금 난다. 서버가 뜨고 내릴 때마다 새 그룹이 생기니 브로커에는 쓰지 않는 그룹 아이디가 계속 쌓인다. 인스턴스 식별자로 그룹 이름을 직접 정해주면 이름이 안정되고 필터링도 이름 규칙 하나로 끝난다. 그 길을 검토했는지 원문은 말하지 않는다.

비즈니스 경로에서 없앤 이중 쓰기가 분석 경로에는 그대로 있다

분석서버는 원본 배달 토픽을 받아 전처리하고 분석 토픽에 재발행한다. 토픽과 서버를 분리해 리소스와 장애 영향범위를 나눈 판단은 명확하다.

배달통합이벤트를 만드는 방식이 흥미롭다. 배달생성 이벤트를 받으면 Redis에 주요 정보를 저장하고, 이후 이벤트가 올 때마다 갱신하고, 배달이 끝나면 Redis에서 지우면서 통합이벤트를 분석 토픽에 발행한다. 카프카에서 읽어 Redis에 쓰고 오프셋을 커밋하는 구조다. 두 시스템에 쓴다.

비즈니스 경로에서 outbox와 CDC까지 동원해 없앤 것이 바로 이 이중 쓰기다. 분석 경로에는 그대로 있다. Redis를 갱신한 뒤 오프셋 커밋 전에 죽으면 같은 갱신이 두 번 일어나고 순서를 바꾸면 갱신이 사라진다. 배달 한 건의 요약 정보가 어긋나는 정도이고 원본 배달 토픽은 그대로 남아 있으니 필요하면 다시 만들 수 있다. 감수할 만한 선택이고 실제로 맞는 선택일 가능성이 높다. 다만 원문은 그것을 선택이라고 말하지 않는다. 한쪽에서는 정합성을 위해 패턴 하나를 통째로 들여왔고 다른 쪽에서는 같은 문제를 그냥 두었는데, 그 온도 차가 의도된 것인지 글에서는 읽히지 않는다.

같은 팀이 Kafka Streams를 이미 운영하고 있다는 점도 걸린다. Streams의 상태저장소는 체인지로그 토픽으로 복구되고 오프셋 커밋과 상태 갱신이 함께 관리된다. 통합이벤트 집계야말로 그 도구에 잘 맞는 모양인데 여기만 Redis다. 통합이벤트 로직이 Streams 도입보다 먼저 만들어졌거나, 다른 서비스에서 조회하기에 Redis가 편했으리라 짐작한다. 추측이다.

정리 정책은 여기서도 빈칸이다. 원문은 “완료된 배달은 Redis에서 삭제하고”라고 적었다. 완료되지 않는 배달은 언제 지워지는지가 없다. 취소된 배달, 어떤 이유로 종료 이벤트가 끝내 오지 않는 배달, 전처리 중 예외로 흘려버린 배달이 남는다. TTL이 걸려 있는지도 나와 있지 않다. 하루 100만 건 규모에서 이 빈칸은 곧 메모리다.

상태저장소는 지우는 이야기가 없으면 미완성이다

실시간 집계는 분석용 배달토픽에서 배달식별자를 키로 한 최신 배달 상태저장소 latest-delivery를 만들고 거기서 배달상태를 키로 한 count-per-status를 만들어 그라파나 게이지로 보여준다. 구조는 깔끔하다.

같은 질문이 여기서도 나온다. latest-delivery는 배달식별자마다 하나의 엔트리를 갖는다. 완료된 배달을 tombstone으로 지우지 않으면 하루 100만 건씩 영구히 쌓인다. RocksDB로 디스크에 앉으니 힙이 바로 터지지는 않겠지만 복구 시간과 체인지로그 크기가 계속 자란다. 원문에는 이 이야기가 없다.

count-per-status는 조금 다른 종류의 함정을 안고 있다. 배달상태가 A에서 B로 바뀌면 집계는 A를 하나 빼고 B를 하나 더해야 맞는다. 스트림을 그냥 상태별로 세면 뺄셈이 없어 숫자가 한 방향으로만 자란다. KTable을 groupBy해서 세면 Streams가 이전 값에 대한 취소 레코드를 함께 흘려보내 그 뺄셈을 대신해준다. 원문이 latest-delivery를 먼저 만들고 거기서 상태별 집계로 넘어가는 순서로 설명하는 걸 보면 이 구조를 골랐을 것으로 보이는데, 코드가 없으니 단정하지는 않겠다. 골랐다면 맞게 간 것이다.

“레코드의 시간을 기준값으로 최신 배달을 판단하며”라는 대목은 첫 절의 이야기와 이어진다. 비즈니스 경로에서는 키와 파티션으로 순서를 보존했다. 그런데 분석 토픽은 원본 토픽을 전처리해 재발행한 결과물이라 여기까지 오는 동안 순서가 한 번 다른 손을 거쳤다. 재발행 시점의 파티셔닝과 처리 지연이 순서를 흐트러뜨릴 수 있으니 도착 순서를 믿지 못하고 레코드 시간을 봐야 한다. 순서를 보존하는 사슬이 한 군데서 끊기면 그 뒤로는 시간 비교라는 비용을 계속 치른다.

해결

세 갈래로 정리된다.

비즈니스 경로는 순서 보존과 at-least-once를 약속한다. MySQL 커밋 순서가 진실이고 outbox 테이블과 Debezium 커넥터와 파티션 키가 그 순서를 컨슈머까지 그대로 나른다. 처리량은 outbox 테이블과 커넥터를 N개로 나눠 확보했고 그 대가로 재분배 절차와 binlog 중복 전송을 떠안았다.

이벤트 버스는 알림만 약속한다. 무엇이 바뀌었는지 말하지 않고 다시 읽으라고만 한다. 순서도 유실도 문제가 되지 않도록 설계된 것이지 순서와 유실을 막아서 안전한 것이 아니다.

분석 경로는 best-effort다. Redis 임시저장소와 오프셋 커밋 사이에 원자성이 없고 S3 싱크 커넥터로 영구 저장한 뒤 Athena로 조회해 사업과 운영 부서가 월단위 분석과 정산에 쓴다. Kafka Streams 집계는 그라파나 대시보드로 나간다.

원문은 스스로 “얇고 넓게 소개하고”라고 밝힌 글이고 실제로 그렇다. 각 구성요소가 무엇인지는 설명하지만 각 구성요소가 무엇을 약속하고 무엇을 약속하지 않는지는 다루지 않는다. 그 한 줄씩만 붙었어도 세 경로의 성격 차이가 훨씬 잘 보였을 것이다.

결과

수치가 없다. 하루 100만 건 이상이라는 규모 말고는 지연도, 처리량도, 커넥터 랙도, 장애 사례도 나오지 않는다. 소개를 목적으로 한 글이니 그럴 수 있다.

그래도 한 곳은 아깝다. outbox 테이블을 N개로 나눴다는 결정은 어딘가에서 랙을 봤다는 뜻이다. 커넥터 하나가 감당하던 초당 몇 건에서 어느 지점부터 밀리기 시작했고 N을 몇으로 늘렸을 때 어떻게 됐는지, 그 한 줄이 이 글에서 가장 값진 문장이 됐을 것이다.

인상 깊었던 점

이 팀 설계의 중심은 순서를 정하는 지점을 하나로 모은 데 있다. 카프카가 순서를 준 것이 아니라, DB가 정한 순서를 카프카가 흐트러뜨리지 않은 것이다. 그래서 “카프카를 쓰니 순서가 보장된다”는 문장은 절반만 맞다. 나머지 절반은 그 순서를 어디서 정했느냐에 있고 그 답이 명확한 시스템과 그렇지 않은 시스템의 운명은 꽤 다르다.

이 블로그에 쓴 Kafka를 걷어내고 DB 작업 큐로: 규모에 맞는 복잡도 되찾기의 outbox와 비교해볼 만하다. 그쪽은 relay가 1초마다 테이블을 폴링해 미발행 행을 읽고, send().get(10s)로 ack를 기다린 뒤 markSent 하는 구조였다. 순서를 relay가 직접 정해야 했고 ack를 기다리는 코드가 필요했다. CDC는 그 두 가지를 binlog에서 공짜로 받는다. 대신 커넥터라는 운영 대상이 하나 늘고 앞에서 본 것처럼 커넥터를 늘릴 때마다 DB가 값을 치른다. 둘 다 outbox라는 같은 이름을 쓰는데 순서를 누가 정하느냐가 다르고 그 차이가 나머지 설계를 거의 다 결정한다.

세 경로가 서로 다른 강도의 약속을 하고 있고 각각 그 자리에 맞는다. 아쉬운 것은 그 차이를 글이 이름 붙여 말하지 않았다는 것뿐이다. 저자는 “팀에서 사용하는 시스템 설계나 구현을 주도해본 경험은 아직 없습니다”라고 적었는데, 2년차가 팀의 시스템 전체를 이만큼 정리해낸 것 자체가 드물다. 남은 빈칸들, 그러니까 컨슈머 멱등성과 outbox 정리 정책과 N을 정한 근거가 그대로 다음 글감이다. 소개 글에서 답을 요구할 것은 아니지만, 답이 있는 팀이라면 그쪽이 훨씬 재미있는 글이 된다.