프로그래밍 및 소프트웨어 개발

Go를 사용해 세션 순서를 유지하는 Kafka 처리 파이프라인 구축 방법

Joshua Oluikpe는 각 세션 내에서 Kafka 메시지를 엄격한 순서로 처리하면서 독립적인 세션은 병렬로 작업할 수 있도록 하는 실용적인 설계를 설명한다. 이 설계는 일관성 해싱, 세션별 Go 워커, 동일 경로 내 재시도, 인접 커밋 포인트 관리를 기반으로 하여 장애 이후 메시지 손실을 방지한다.

2026-10-07
4 분 읽기
1 조회수
certi.news Editorial Team
Go를 사용해 세션 순서를 유지하는 Kafka 처리 파이프라인 구축 방법

수천 개의 세션이 동일한 파티션을 공유할 때 Kafka의 파티션 내부 순서에만 의존하는 것으로는 충분하지 않다. 대화 및 순차 작업 시스템에서는 «목적지를 파리로 설정해»와 같은 수정 메시지가 «런던행 항공편을 예약해»라는 메시지 뒤에 도착할 수 있지만, 이를 첫 번째 메시지보다 먼저 실행하면 맥락이 깨진다. InfoQ에서 Joshua Oluikpe가 설명한 해결책은 순서 보장의 일부 책임을 Go로 작성된 애플리케이션 계층으로 옮긴다.

세션 내 순서와 세션 간 병렬성

이 설계는 두 단계의 워커를 사용한다. 분배 계층은 Kafka 레코드를 수신하고 일관성 해싱을 사용해 세션 ID에 따라 전달한다. 따라서 동일한 세션의 메시지는 동일한 워커로 전달되는 반면, 서로 다른 세션은 여러 워커에 분산될 수 있다. 이후 각 활성 세션에는 전용 goroutine이 할당되어 한 번에 하나의 메시지를 도착 순서대로 처리한다.

이렇게 하면 하나의 느린 세션이 다른 세션을 가로막지 않는다. 평면적인 워커 풀에서는 한 세션의 재시도가 워커 전체를 점유할 수 있다. goroutine은 필요할 때 생성되고 일정 시간 비활성 상태가 지속되면 제거되므로, 리소스 사용량은 가능한 전체 세션 수가 아니라 활성 세션 수에 연동된다.

메시지를 건너뛰지 않는 재시도

재시도는 동일한 세션의 goroutine 내부에서 수행된다. 종속 서비스의 시간 초과처럼 일시적인 오류가 발생하면 워커는 일정한 무작위성을 포함한 지수 백오프를 사용해 대기한 뒤, 후속 메시지로 넘어가기 전에 동일한 메시지를 다시 처리한다. 다음 메시지들은 실패한 메시지 뒤의 채널에 남아 있으므로 네 번째 메시지가 세 번째 메시지를 앞지를 수 없다.

반면 잘못된 데이터와 같은 복구 불가능한 오류는 직접 배달 불가 메시지 목록(DLQ)으로 이동한다. 일시적인 오류와 최종 오류를 이렇게 구분하면 재시도 관리를 위한 별도의 상태 머신이 필요하지 않지만, 백오프 기간 동안 특정 세션의 메시지가 누적될 수 있다는 의미이기도 하다.

안전한 커밋과 장애 이후 복구

여러 세션이 메시지를 병렬로 처리할 때는 완료된 가장 높은 offset만 커밋해서는 안 된다. 예를 들어 offset 104가 102보다 먼저 완료된 뒤 104를 커밋하면, 소비자 장애로 인해 102가 영구적으로 건너뛰어질 수 있다. 따라서 시스템은 각 파티션별로 처리 중인 offset과 완료된 offset을 추적하며, 인접한 형태로 완료된 가장 높은 offset까지만 커밋 포인트를 이동한다.

그 결과 재시작은 더 이른 offset에서 시작되어 일부 메시지를 다시 처리할 수 있다. 이 설계는 중복 제거를 지원하기 위해 고정된 이벤트 ID를 사용하지만, «정확히 한 번» 실행을 달성한다고 주장하지는 않는다. 외부 효과는 여전히 종속 서비스가 중복을 멱등적으로 처리할 수 있는 능력에 좌우된다. 드문 경우에는 파티션의 진행을 복구하기 위해 일정한 시간 제한이 지난 후 정체된 간극을 건너뛸 수 있으며, 이를 운영 실패로 기록하고 메시지를 DLQ로 보낸다. 이는 연속성과 처리 완전성 사이의 명시적인 절충이다.

실제 운영 환경에 필요한 것

순서 보장은 기본 로직만으로 완성되지 않는다. 파티션 재분배 중에는 소비자가 새 레코드 수신을 중지하고, 진행 중인 작업이 끝나기를 기다린 뒤 소유권을 넘기기 전에 마지막 인접 watermark를 커밋한다. 워커 채널이 가득 차면 메시지를 버리는 대신 일시적으로 가져오기를 중지하며, 이미 가져온 레코드는 각 파티션별 저장소에 보관해 더 최신 메시지가 앞서 진행되지 않도록 한다.

또한 시스템은 정체된 offset을 개별적으로 탐지하고, 비동기 Kafka 경계를 가로지르는 추적을 연결하며, 가져오기를 중지하는 동안 자동 확장에 필요한 지연 신호를 유지하고, 종료 전에 작업을 비우고 DLQ를 기록해야 한다. 제품은 비동기 생산 과정에서도 메시지 순서를 보장하기 위해 전송 측에도 일관성 해싱을 적용한다.

결과와 한계

50,000개의 메시지와 메시지의 10%에 일시적인 오류를 주입한 10개 세션으로 구성된 합성 테스트에서 시스템은 순서 위반 없이 초당 14,027개의 메시지를 기록했으며, 중간 지연 시간은 11.7밀리초였다. 1,000개 세션을 대상으로 한 부하 테스트에서는 전송 오류 없이 목표치 50,000개 중 초당 48,805개의 실효 전송률을 기록했지만, 꼬리 지연 시간은 200밀리초 예산을 초과했다. 이는 세션 순서 메커니즘이 아니라 API Gateway의 메시지 속도 할당량 때문이었다.

작성자는 이 시스템이 눈에 띄는 순서 위반 없이 4,000만 개가 넘는 운영 메시지를 처리했으며, 그중 460개 메시지만 DLQ로 전달되었다고 말한다. 이는 약 0.002%에 해당한다. 그러나 이 수치는 운영 모니터링에 기반한 것이며 모든 메시지 시퀀스를 포괄적으로 검사한 결과는 아니다. 또한 더 높은 부하 테스트에서는 세션별 순서를 독립적으로 검증하지 않았다. 따라서 이 설계의 실질적인 가치는 순서를 Kafka 수준의 가정에서 애플리케이션이 강제하는 속성으로 전환하는 데 있으며, 동시에 명시적인 시퀀스 테스트와 종속 서비스의 멱등성 보장이 여전히 필요하다.

뉴스 출처
InfoQ - Architecture Articles
원문 보기 ↗
c
작성자

certi.news Editorial Team

같은 카테고리

추천 기사

모든 뉴스 보기