전체 그래프
Kafka

Kafka Consumer

data-engineeringtool/kafkaconcept/consumer

상위: Kafka

요약

Consumer는 Kafka Topic에서 데이터를 읽는 역할을 합니다. "구독(subscribe)" 구조로 동작하며, Consumer Group을 통해 병렬 처리가 가능합니다.

Consumer 동작 방식

image 82 1.png

Polling 방식

  • 브로커가 푸시하는 것이 아님
  • Consumer가 직접 요청해서 가져감
  • Pull 기반 모델

image 84 1.png

멀티 컨슈밍

  • 하나의 토픽에 대해 여러 Consumer가 병렬로 처리 가능
  • Consumer Group을 통해 관리

image 85 1.png

Consumer Lag

image 83 1.png

  • Consumer Lag: 최신 메시지의 offset - 현재 읽고 있는 offset
  • record-lag-max: 파티션 중 가장 높은 Lag 값
  • Lag이 증가하면 처리 지연 의미

Consumer Group

image 86 1.png

  • 특정 토픽에 대해 접근 권한을 가진 Consumer 집단
  • 같은 Group 내 Consumer들은 파티션을 나눠서 처리
  • 다른 Group은 독립적으로 메시지 소비 가능

Fetch & Commit

Fetch (패치)

image 87 1.png

  • Consumer가 브로커에 데이터를 요청
  • 일정량의 데이터를 한 번에 가져옴

Commit (커밋)

image 88 1.png

  • 특정 offset까지 읽었음을 선언
  • 재시작 시 마지막 커밋 위치부터 읽음

Group Coordinator

image 89 1.png

image 90 1.png

  • Coordinator 역할의 Broker가 Consumer Group을 관리
  • Heartbeat: Consumer가 살아있는지 주기적으로 확인
  • Rebalancing: Consumer 수/파티션 변화 시 파티션 재분배

Rebalancing

발생 조건

  1. 새로운 Consumer가 추가될 경우
  2. 기존 Consumer가 제거될 경우
  3. 파티션 수 변경

Rebalancing 과정

image 91 1.png

image 92 1.png

  1. Coordinator가 모든 Consumer의 파티션 소유권 회수, 일시정지
  2. JoinGroup 요청 대기 후 가장 먼저 응답한 Consumer를 리더로 지정
  3. 리더가 새로운 파티션 할당 결과 계산 → Coordinator에 전달 → 각 Consumer에 전파

파티셔닝 전략

image 93 1.png

1. RangeAssignor

토픽 기준으로 순서대로 파티션 분배

image 94 1.png

2. RoundRobinAssignor

전체 파티션을 골고루 할당 (1개씩 분배)

image 95 1.png

3. StickyAssignor

기존 할당 정보를 고려해 변경 최소화

참고: partition.assignment.strategy의 실제 기본값은 [RangeAssignor, CooperativeStickyAssignor] 리스트다 (Kafka 2.4+). 리스트 첫 항목인 RangeAssignor가 우선 적용되며, RangeAssignor를 빼는 롤링 재시작 한 번으로 CooperativeStickyAssignor(점진적 협력 리밸런싱)로 무중단 전환할 수 있게 설계된 값. StickyAssignor 자체는 기본값이 아니라 선택 옵션이다.

image 96 1.png

→ 리밸런싱 시간 단축 + 캐시 효율성 유지

Commit 전략

자동 커밋

  • enable_auto_commit=True
  • 주기적으로 자동 커밋 (기본값: 5초)
  • 실패 시 유실 가능

수동 커밋

  • enable_auto_commit=False
  • 정상 처리 후 명시적으로 커밋
  • 안전하지만 부하 증가

권장: 배치 단위로 마지막 메시지만 커밋 (안전성과 효율의 균형)

Transaction

Kafka Consumer

  • Producer ↔ Consumer 간 정확히 한 번(Exactly Once) 처리 보장
  • 트랜잭션 단위로 메시지를 묶고 커밋
  • 커밋되지 않으면 Consumer는 수신 안 함
  • 실패 시 전체 트랜잭션 무효화

Python 예제

from kafka import KafkaConsumer

consumer = KafkaConsumer(
    'test-topic',
    bootstrap_servers='localhost:9092',
    auto_offset_reset='earliest',
    enable_auto_commit=True,
    group_id='my-group',
    value_deserializer=lambda v: v.decode('utf-8')
)

# 메시지 수신
print("Listening for messages...")
for message in consumer:
    print(f"Received: {message.value}")

주요 설정

설정설명기본값
group.idConsumer Group ID없음
auto.offset.reset초기 offset 설정latest
enable.auto.commit자동 커밋 여부True
max.poll.records한 번에 가져올 최대 메시지 수500
session.timeout.msHeartbeat 타임아웃10000ms