LLM WikiAccess-protected knowledge portal

WIKI

프로듀서와 컨슈머 설계

데이터 전달 보장 수준 Kafka를 실제로 쓸 때 가장 먼저 결정해야 하는 것은 "메시지를 얼마나 안전하게 보낼 것인가"다. 이는 acks 설정 과 컨슈머 오프셋 커밋 전략 두 가지로 결정된다. 전달 보장은 세 단계로 나뉜다. 보장 수준 의미 적합한 상황 At most once 재전송 없음. 유실 가능 로그처럼 일부 손실이 허용될 때 At least once 재전송으로 중복 가능 처리가 idempotent할 때 기본값 Exac

경로human/study/content/kafka/02-producer-consumer.md
카테고리Study
태그#consumer #content #kafka #producer #study

데이터 전달 보장 수준

Kafka를 실제로 쓸 때 가장 먼저 결정해야 하는 것은 "메시지를 얼마나 안전하게 보낼 것인가"다. 이는 acks 설정컨슈머 오프셋 커밋 전략 두 가지로 결정된다.

전달 보장은 세 단계로 나뉜다.

보장 수준의미적합한 상황
At-most-once재전송 없음. 유실 가능로그처럼 일부 손실이 허용될 때
At-least-once재전송으로 중복 가능처리가 idempotent할 때 (기본값)
Exactly-once중복 없이 정확히 한 번결제, 재고 등 중복이 치명적일 때

Producer: 메시지를 브로커에 쓰기

acks 설정

프로듀서가 produce() 요청을 보내면 브로커는 설정한 acks 수준에 따라 응답 시점이 달라진다.

acks=all과 함께 min.insync.replicas=2를 설정하면, "적어도 두 브로커에 기록됐을 때만 성공으로 본다"는 보장이 생긴다.

멱등 프로듀서 (Idempotent Producer)

네트워크 오류 후 재시도할 때 같은 메시지가 두 번 브로커에 기록될 수 있다. 이를 막는 것이 멱등 프로듀서다.

Kafka 3.0부터 enable.idempotence=true가 기본값이다. 브로커는 프로듀서에게 Producer ID(PID)를 부여하고, 파티션별 시퀀스 번호를 추적한다. 재전송된 메시지의 PID+시퀀스가 이미 기록된 것과 같으면 중복으로 처리해 무시한다.

중요한 한계: 멱등성은 단일 프로듀서 세션 안에서만 보장된다. 프로듀서가 재시작되면 PID가 바뀌어 중복을 감지할 수 없다. 프로세스 재시작을 넘는 exactly-once가 필요하다면 트랜잭션을 써야 한다.

파티셔닝 전략

메시지가 어느 파티션으로 갈지는 파티셔너가 결정한다.

파티션 수를 늘리면 기존 키-파티션 매핑이 깨진다. 순서가 중요하다면 파티션 수를 초기에 넉넉하게 잡는다.

Consumer: 메시지를 브로커에서 읽기

오프셋 커밋 전략

컨슈머는 메시지를 읽은 후 처리 완료를 __consumer_offsets오프셋 커밋으로 기록한다. 재시작이나 리밸런싱 후 이 커밋 지점부터 이어서 읽는다.

자동 커밋 (enable.auto.commit=true, 기본값): auto.commit.interval.ms마다 마지막 poll 오프셋을 자동으로 커밋한다. 간단하지만, 커밋된 오프셋 이후 처리 도중 컨슈머가 죽으면 다음 재시작에서 그 구간을 건너뛴다(at-most-once 위험).

수동 커밋: 처리가 완료된 뒤 직접 commitSync() 또는 commitAsync()를 호출한다. 처리 전 커밋하면 유실, 처리 후 커밋하면 중복 위험이 각각 남는다. 일반적으로 처리 후 커밋을 권장한다(at-least-once).

while (true) {
    records = consumer.poll(timeout)
    for record in records:
        process(record)          // 먼저 처리
    consumer.commitSync()        // 처리 후 커밋
}

exactly-once가 필요하면 오프셋을 외부 저장소(DB 등)에 처리 결과와 원자적으로 함께 쓴다.

Consumer Group 리밸런싱

컨슈머 그룹에 새 인스턴스가 합류하거나 기존 인스턴스가 떠나면 리밸런싱(rebalancing) 이 일어나 파티션이 재분배된다.

Eager 리밸런싱 (기본): 모든 컨슈머가 파티션을 일제히 반납하고 재할당. 리밸런싱 동안 전체 그룹이 소비를 멈춘다(stop-the-world).

Cooperative Sticky 리밸런싱 (partition.assignment.strategy=CooperativeStickyAssignor): 실제로 이동이 필요한 파티션만 점진적으로 재할당. 이동하지 않는 파티션은 계속 소비할 수 있어 처리 중단 시간을 최소화한다. Kafka 2.4+, 클라이언트 2.4+ 이상에서 사용 가능.

흐름 요약

Producer acks=all
멱등·트랜잭션
파티셔너
key hash / sticky
Broker (Leader) 로컬 로그 쓰기 ISR 복제 대기 ACK 응답
Consumer poll() 처리 commitSync/
Async
__consumer_offsets
커밋 기록
Group Coordinator
멤버십·리밸런싱
프로듀서 → 브로커 → 컨슈머의 전달 흐름과 보장 레이어

Producer 주요 설정 요약

설정기본값 (Kafka 3.0+)역할
acksall쓰기 내구성
enable.idempotencetrue재시도 중복 방지
retriesMAX_INT재시도 횟수
max.in.flight.requests.per.connection5순서 보장과 처리량의 균형
compression.typenonesnappy/lz4/zstd로 대역폭 절감 가능
linger.ms0배치 지연 허용 시간
batch.size16384 bytes배치 최대 크기

Consumer 주요 설정 요약

설정기본값역할
enable.auto.committrue자동 커밋 여부
auto.commit.interval.ms5000자동 커밋 주기
auto.offset.resetlatest커밋 없을 때 시작 위치 (earliest/latest)
max.poll.records500poll 1회 최대 레코드 수
session.timeout.ms45000하트비트 없을 때 컨슈머 제거 기준
partition.assignment.strategyRangeAssignor리밸런싱 전략

Open Questions

References