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

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

멀티 컨슈밍
- 하나의 토픽에 대해 여러 Consumer가 병렬로 처리 가능
- Consumer Group을 통해 관리

Consumer Lag

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

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

- Consumer가 브로커에 데이터를 요청
- 일정량의 데이터를 한 번에 가져옴
Commit (커밋)

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


- Coordinator 역할의 Broker가 Consumer Group을 관리
- Heartbeat: Consumer가 살아있는지 주기적으로 확인
- Rebalancing: Consumer 수/파티션 변화 시 파티션 재분배
Rebalancing
발생 조건
- 새로운 Consumer가 추가될 경우
- 기존 Consumer가 제거될 경우
- 파티션 수 변경
Rebalancing 과정


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

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

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

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

→ 리밸런싱 시간 단축 + 캐시 효율성 유지
Commit 전략
자동 커밋
enable_auto_commit=True- 주기적으로 자동 커밋 (기본값: 5초)
- 실패 시 유실 가능
수동 커밋
enable_auto_commit=False- 정상 처리 후 명시적으로 커밋
- 안전하지만 부하 증가
권장: 배치 단위로 마지막 메시지만 커밋 (안전성과 효율의 균형)
Transaction

- 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.id | Consumer Group ID | 없음 |
auto.offset.reset | 초기 offset 설정 | latest |
enable.auto.commit | 자동 커밋 여부 | True |
max.poll.records | 한 번에 가져올 최대 메시지 수 | 500 |
session.timeout.ms | Heartbeat 타임아웃 | 10000ms |