「exactly-once 켰으니 중복 없음」 PR 한 줄: Kafka 보장이 끝나는 경계와 면접 답

팀그릿·어제
post-thumbnail

금요일 오후에 코드 리뷰 알림이 옵니다. PR 제목은 「정산 컨슈머 exactly-once 적용」입니다. 설명란에는 「프로듀서 멱등성을 켜고 Kafka 트랜잭션으로 묶었습니다. 이제 결제 이벤트가 두 번 처리될 일은 없습니다.」라고 적혀 있습니다. diff 를 보면 컨슈머는 이벤트를 읽고, 정산 테이블에 쓰고, 결제사 API 를 호출합니다. 리뷰어가 댓글을 한 줄 남깁니다. 「결제사 호출도 그 트랜잭션 안에 들어가나요?」

이 댓글 하나로 「Kafka 는 exactly-once 를 보장하나요?」의 답이 갈립니다. 보장은 합니다. 다만 범위는 Kafka 안의 읽기, 처리, 쓰기까지입니다. DB 와 외부 API 는 그 범위 밖입니다. 그릿 딥다이브 Vol.2 · Kafka 3주차 「생존」에서 다룬 내용입니다.

exactly-once 는 메시지가 딱 한 번 간다는 뜻인가요?

아닙니다. 전송과 재시도는 여러 번 일어날 수 있습니다. exactly-once 는 그래도 처리의 효과가 한 번만 반영된다는 뜻입니다. 그래서 보장의 범위는 효과가 남는 저장소가 어디까지인지로 정해집니다.

Kafka 트랜잭션은 두 가지를 묶습니다. Kafka 토픽에 쓰는 일과 Kafka 소비 오프셋(컨슈머 그룹이 어디까지 읽었는지 기록한 위치)입니다. MySQL INSERT, 이메일 발송, 결제사 HTTP 요청은 같은 트랜잭션에 자동으로 들어오지 않습니다. 그래서 리뷰어 질문의 답은 「들어가지 않습니다」입니다. 결제 API 가 성공하고 커밋 전에 프로세스가 멈추면 재처리 때 API 를 다시 호출할 수 있습니다. 이미 실행한 외부 HTTP 요청을 Kafka 트랜잭션이 롤백하지 못합니다.

중복은 어디서 생기나요?

중복은 두 종류로 나뉩니다. 하나는 Kafka 로그에 같은 업무 이벤트가 두 번 기록되는 경우입니다. 다른 하나는 로그의 레코드 하나를 컨슈머가 두 번 처리하는 경우입니다.

로그 중복은 프로듀서 재시도에서 생깁니다. 브로커가 레코드를 기록했는데 응답이 끊기면 프로듀서는 실패로 판단합니다. 이때 다시 보내면 레코드가 두 번 남을 수 있습니다. 재시도를 끄면 중복은 줄지만 일시적인 네트워크 오류나 리더 변경 때 유실될 가능성이 커집니다.

재처리 중복은 처리와 오프셋 커밋의 순서에서 생깁니다. Kafka 오프셋 커밋과 외부 DB 트랜잭션은 기본적으로 하나의 원자 연산이 아닙니다. 어느 쪽을 먼저 해도 두 연산 사이에서 멈출 수 있는 구간이 남습니다.

순서사이에서 멈추면이름
처리 후 커밋다음 인스턴스가 같은 레코드를 다시 처리at-least-once
커밋 후 처리다음 인스턴스가 그 레코드를 건너뜀at-most-once

중복을 없애려고 커밋을 처리 앞으로 옮기면 중복이 유실로 바뀔 뿐입니다. 그래서 처리 후 커밋에 멱등한 업무 로직을 함께 쓰는 경우가 많습니다.

멱등 프로듀서는 무엇을 막나요?

멱등 프로듀서는 배치마다 producer ID, producer epoch(같은 ID를 쓰는 프로듀서의 세대 번호), 파티션별 sequence 를 붙입니다. 브로커는 프로듀서마다 최근 sequence 를 추적합니다. 같은 배치가 재전송되면 브로커가 알아보고 로그에 다시 쓰지 않습니다.

sequence 는 레코드 내용의 해시가 아닙니다. 같은 프로듀서 세션이 같은 파티션에 보낸 레코드마다 하나씩 올라가는 번호입니다. producer ID 도 애플리케이션이 정하지 않습니다. 프로듀서가 InitProducerId 요청을 보내면 브로커가 발급합니다. 프로세스가 재시작하면 새 ID를 받습니다.

Kafka 3.9.0 에서는 enable.idempotence=true 가 기본입니다. 다음 조건이 함께 필요합니다.

  • acks=all
  • retries > 0
  • max.in.flight.requests.per.connection <= 5

숫자 5에는 이유가 있습니다. 브로커는 파티션마다 producer ID 별로 최근 다섯 개 배치의 sequence 를 기억합니다. 동시에 떠 있는 요청이 이보다 많으면 재전송된 배치를 판정할 근거가 사라집니다. 기대한 sequence 와 어긋난 배치를 브로커가 거부하면 프로듀서는 OutOfOrderSequenceException 을 받습니다.

설정 충돌도 확인해야 합니다. enable.idempotence=true 를 명시하고 충돌하는 값을 넣으면 ConfigException 이 발생합니다. 멱등 설정을 명시하지 않고 acks=1 이나 retries=0 을 넣으면 멱등이 꺼질 수 있습니다. 오래된 설정 한 줄 때문에 기본 보호가 꺼질 수 있으므로 시작 로그에서 최종 적용 설정을 확인합니다.

멱등 프로듀서가 막지 못하는 것도 분명합니다.

  • 애플리케이션이 예외를 잡아 새로 호출한 send() 는 새 레코드 전송입니다. 새 sequence 를 받으므로 브로커는 중복으로 보지 않습니다.
  • 다른 프로듀서 인스턴스가 같은 업무 이벤트를 보낸 경우도 막지 못합니다.
  • sequence 는 파티션마다 따로 매깁니다. 그래서 파티션 하나 안의 중복만 막습니다.
  • 컨슈머가 같은 레코드를 두 번 처리하는 문제와는 관계가 없습니다.

함께 보면 좋은 글: Redis 백업 중 메모리가 두 배로 찍혔다: 서버를 늘리기 전에, 면접에서 답하기 전에 볼 세 가지

Kafka 트랜잭션은 어디까지 묶나요?

Kafka 트랜잭션은 여러 토픽 파티션에 쓰는 일과 입력 오프셋 커밋을 한 트랜잭션으로 묶습니다. 입력 토픽 A를 읽어 출력 토픽 B에 쓰는 흐름이라면 출력만 남거나 오프셋만 전진하는 상태가 생기지 않습니다. 둘 다 커밋되거나 둘 다 중단됩니다.

프로듀서에 안정적인 transactional.id 를 설정하고 아래 순서로 호출합니다. 예외 처리는 확인하기 쉽도록 뺐습니다.

producer.initTransactions();
while (running) {
    ConsumerRecords<String, Event> records = consumer.poll(timeout);
    producer.beginTransaction();
    for (ConsumerRecord<String, Event> record : records) {
        producer.send(toOutputRecord(record));
    }
    producer.sendOffsetsToTransaction(
            offsetsOf(records), consumer.groupMetadata());
    producer.commitTransaction();
}

예외는 두 가지로 나눠 다룹니다. ProducerFencedException 처럼 이 인스턴스가 이미 밀려난 경우는 재시도해도 회복되지 않습니다. 이때는 프로듀서를 닫고 인스턴스를 내립니다. 그 밖의 KafkaException 은 abortTransaction() 으로 트랜잭션을 중단합니다. 오프셋 커밋도 함께 취소되므로 다음 poll 에서 같은 구간을 다시 가져옵니다.

transactional.id 가 재시작 전후로 같아야 하는 이유가 있습니다. 같은 ID로 새 인스턴스가 initTransactions() 를 호출하면 트랜잭션 코디네이터가 그 ID의 epoch 를 올립니다. 그 뒤로 낡은 epoch 를 단 이전 인스턴스가 쓰려고 하면 ProducerFencedException 을 받습니다. 재시작할 때마다 새 UUID 를 쓰면 코디네이터가 이전 세대를 알아보지 못합니다. 그러면 두 세대가 서로 다른 ID로 동시에 쓰고, 좀비 차단(밀려난 이전 인스턴스의 쓰기를 막는 장치)이 동작하지 않습니다.

컨슈머 쪽에서는 consumer.groupMetadata() 가 같은 역할을 합니다. 이 값에는 그룹 ID, 현재 generation(리밸런싱마다 하나씩 오르는 세대 번호), 멤버 ID가 들어 있습니다. 그룹 코디네이터는 낡은 generation 의 오프셋 커밋을 거부합니다. Kafka Streams 의 processing.guarantee=exactly_once_v2 가 이 방식(KIP-447)을 씁니다.

읽는 쪽 설정도 바꿔야 합니다. 중단된 트랜잭션의 레코드도 로그에는 물리적으로 남습니다. 커밋과 중단 표시는 제어 레코드로 같은 파티션에 기록됩니다. isolation.level=read_committed 컨슈머는 이 표시를 읽고 중단된 레코드를 건너뜁니다. 기본값 read_uncommitted 는 중단된 레코드도 읽습니다. 그리고 트랜잭션이 오래 열려 있으면 read_committed 컨슈머는 LSO(Last Stable Offset, 결과가 정해지지 않은 트랜잭션 바로 앞 위치)에서 멈춥니다. 소비가 밀린 것처럼 lag 이 쌓여 보입니다.

결제 API 호출은 누가 막나요?

Kafka 밖의 효과는 끝점(효과가 실제로 남는 DB나 외부 API)마다 따로 멱등하게 만들어야 합니다.

  • DB에는 업무 ID 고유 제약이나 멱등 소비 기록을 둡니다.
  • 외부 API에는 업무 ID로 만든 idempotency key(같은 요청의 재호출임을 알리는 식별자)를 보냅니다.
  • 처리 상태를 단계별로 저장하고 재시작하면 멈춘 단계부터 이어서 수행합니다.
  • DB 변경과 발행할 이벤트를 함께 기록하는 transactional outbox 를 검토합니다.

멱등 소비 기록은 이미 처리한 이벤트를 표시해 두는 테이블입니다. 처리 전에 표시를 넣어 보고, 표시가 이미 있으면 업무 처리를 건너뜁니다. MySQL 이라면 이렇게 확인할 수 있습니다.

CREATE TABLE processed_event (
    event_id     VARCHAR(64) PRIMARY KEY,
    processed_at TIMESTAMP   NOT NULL DEFAULT CURRENT_TIMESTAMP
);

-- 영향받은 행이 0이면 이미 처리한 이벤트이므로 업무 처리를 건너뜁니다.
INSERT IGNORE INTO processed_event (event_id) VALUES (?);

표시와 업무 변경은 같은 DB 트랜잭션 안에 있어야 합니다. 둘을 나누면 표시만 남고 업무 변경은 빠진 상태가 생길 수 있습니다. 그러면 재처리 때 그 이벤트는 처리 완료로 걸러집니다. 중복을 막으려던 장치 때문에 오히려 누락이 생깁니다.

키를 무엇으로 정하는지도 중요합니다. (topic, partition, offset) 을 키로 쓰면 같은 레코드의 재처리는 걸러집니다. 하지만 프로듀서가 두 번 보내 서로 다른 오프셋에 남은 두 레코드는 통과합니다. 업무 ID를 키로 쓰면 프로듀서가 만든 중복까지 같은 키로 모입니다.

그래서 외부 결제 API로 끝나는 파이프라인에서 현실적인 목표는 at-least-once 전달과 끝점 멱등성입니다. 중복이 없다고 가정하지 않습니다. 중복이 와도 결과가 한 번 처리한 것과 같게 만듭니다.

이런 CS 질문 하나를 아침마다 짧게 같이 풀어 보는 오픈채팅방이 있습니다. 개발자: 데일리 CS 역량 강화 챌린지 들어가 보기 →

면접에서 이 질문을 받는다면

「Kafka 는 exactly-once 를 보장하나요?」에는 이렇게 답할 수 있습니다.

Kafka 안의 읽기, 처리, 쓰기까지는 보장합니다. 멱등 프로듀서는 같은 세션의 전송 재시도를 파티션 안에서 걸러 냅니다. 트랜잭션은 출력 토픽 쓰기와 입력 오프셋 커밋을 함께 커밋하거나 함께 중단하고, 컨슈머는 read_committed 로 중단된 레코드를 건너뜁니다. 다만 exactly-once 는 효과가 한 번 반영된다는 뜻이라서 외부 DB나 결제 API에 남는 효과는 이 범위에 들어오지 않습니다. 그래서 외부 연동은 처리 후 커밋하는 at-least-once 소비에 업무 ID 고유 제약과 idempotency key 를 더해 설계합니다.

면접관이 이어서 「멱등 프로듀서가 기본으로 켜져 있는데 왜 업무 중복이 생기나요?」라고 물을 수 있습니다. sequence 는 한 프로듀서 세션이 한 파티션에 보낸 레코드에만 매기는 번호라는 점부터 말합니다. 애플리케이션이 다시 호출한 send() 와 재시작한 프로세스는 새 sequence 나 새 producer ID를 받습니다. 컨슈머가 커밋 전에 멈춰서 생기는 재처리는 프로듀서와 관계없습니다. 여기에 acks=1 같은 오래된 설정이 멱등을 끌 수 있다는 점을 덧붙이면 됩니다.

이런 꼬리 질문을 매일 하나씩 익명으로 같이 풀어 보는 방도 있습니다. 개발자: 데일리 익명 면접 챌린지 들어가 보기 →

같은 장에서 하나 더: 자동 커밋은 중복을 만드나요, 누락을 만드나요?

Kafka 3.9.0 컨슈머는 기본으로 enable.auto.commit=true, auto.commit.interval.ms=5000(5초)입니다. 자동 커밋은 poll() 흐름에서 주기를 확인하고, 반환한 배치 이후 위치를 커밋합니다. 프로세스가 멈추면 최대 한 주기 분량을 다시 읽습니다. 기본값이라면 마지막 5초 동안 처리한 레코드입니다. close() 로 정상 종료하거나 리밸런싱으로 파티션을 반납할 때도 커밋이 일어납니다. kill -9 나 컨테이너 강제 중단은 이 두 경우를 모두 건너뜁니다. 결과가 어느 쪽인지는 처리 구조에 달려 있습니다. 배치를 모두 처리하고 나서 다음 poll() 을 부르면 중복 쪽입니다. 레코드를 작업 스레드에 넘기고 곧바로 poll() 을 부르면 처리보다 커밋이 먼저 나갈 수 있어 누락 쪽입니다. enable.auto.commit=true 한 줄만 보고는 전달 의미를 알 수 없습니다.

같은 장에서 다룬 나머지는 다음과 같습니다.

  • 리밸런싱이 소비 중복을 만드는 방식: session.timeout.ms 와 max.poll.interval.ms 가 서로 다른 것을 감시하는 이유, eager 와 협력적 리밸런싱의 차이, onPartitionsRevoked 와 onPartitionsLost 를 나눠 구현해야 하는 경우를 다룹니다.
  • 처리 실패 시 커밋을 어디까지 전진시킬지: seek 로 되감기, pause 후 재시도, dead letter 토픽, 배치 전체 재처리가 각각 어떤 중복과 누락을 남기는지 비교합니다.
  • 비동기 커밋의 순서 역전: commitAsync() 응답이 늦게 도착해 최신 커밋을 덮는 경로와 콜백에서 이를 막는 방법을 다룹니다.

함께 읽기

이 글의 출처

이 글은 팀그릿 「그릿 딥다이브 Vol.2 · Kafka」 3주차 「생존」에서 뽑았습니다. 리밸런싱, 실패 시 커밋 전략, 비동기 커밋 순서 역전의 풀이와 설정은 책에 있습니다.

그릿 딥다이브 Vol.2 · Kafka 목차 보기 →

profile
면접 꼬리 질문에서 말문이 막혀 본 적 있나요? 저녁마다 그 질문 하나를 익명으로 같이 풀고, 아침마다 CS 질문 하나를 같이 봅니다. 입장 링크는 소개 탭에 있습니다. Redis·Kafka 딥다이브 책을 쓴 팀그릿 · teamgrit.co

0개의 댓글