본문으로 건너뛰기
김신건의 로그

[Distributed] Kafka Consumer Group: rebalancing, offset, lag

· 수정 · 📖 약 2분 · 604자/단어 #kafka #consumer-group #rebalancing #backend #distributed
kafka-consumer-group, Kafka consumer group, Kafka rebalancing, cooperative rebalancing, consumer lag, static membership, partition assignment, Group Coordinator

정의

Consumer Group = 같은 group.id를 가진 consumer들이 함께 한 topic을 분담 소비. partition 단위로 분배.

핵심 특성:

  • 한 partition → 한 consumer (group 안에서)
  • 여러 Consumer Group은 같은 topic을 독립적으로 소비 가능
  • Group Coordinator (Kafka broker)가 멤버십과 partition 할당 관리

Partition 할당

flowchart LR
    subgraph T["Topic: orders (6 partitions)"]
        P0[p0] & P1[p1] & P2[p2] & P3[p3] & P4[p4] & P5[p5]
    end
    subgraph CG["Consumer Group"]
        C1["Consumer 1"] --> P0 & P1
        C2["Consumer 2"] --> P2 & P3
        C3["Consumer 3"] --> P4 & P5
    end

규칙:

  • partition 수 >= consumer 수면 모두 분배
  • partition 수 < consumer 수일부 consumer idle
  • 한 partition → 한 consumer (group 안)

Group Coordinator

Group Coordinator = consumer group의 멤버십과 rebalancing을 관리하는 Kafka broker.

sequenceDiagram
    autonumber
    participant C as "Consumer"
    participant GC as "Group Coordinator (Broker)"
    participant ZK as "__consumer_offsets topic"

    C->>GC: FindCoordinator (group.id)
    GC-->>C: Coordinator broker 주소
    C->>GC: JoinGroup
    GC->>GC: 모든 멤버 대기 (sync barrier)
    GC-->>C: JoinGroup Response (leader 선출)
    Note over C: Group Leader가 partition 할당 계산
    C->>GC: SyncGroup (할당 결과)
    GC-->>C: SyncGroup Response (내 할당)
    C->>ZK: offset commit
  • __consumer_offsets 내부 topic에 offset 저장
  • Group Leader (consumer 중 하나)가 실제 partition 할당 계산
  • Coordinator는 할당 결과를 모든 멤버에게 배포

Rebalancing

consumer 추가 / 제거 / 죽음 → partition 재할당:

sequenceDiagram
    autonumber
    participant C1 as "Consumer 1"
    participant C2 as "Consumer 2 (new)"
    participant Coord as "Group Coordinator"

    C1->>Coord: heartbeat
    C2->>Coord: join group
    Note over Coord: rebalance 시작
    Coord->>C1: stop fetching (revoke)
    C1->>Coord: ack + commit offset
    Coord->>C1: "new assignment (p0, p1, p2)"
    Coord->>C2: "new assignment (p3, p4, p5)"
    C1->>C1: resume fetching
    C2->>C2: start fetching

Eager vs Cooperative

Eager (옛 기본)Cooperative (Kafka 2.4+)
Rebalance 시모든 consumer 정지영향받는 partition만
처리 중단길음 (수십 초 가능)짧음
Throughput하락 큼적음
구현RangeAssignorCooperativeStickyAssignor

IMPORTANT

Cooperative rebalancing2026 권장 기본. partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor.

Cooperative Sticky Assignor 설정

# consumer.properties
group.id=my-consumer-group
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

# 또는 여러 전략 조합 (마이그레이션 시)
partition.assignment.strategy=\
  org.apache.kafka.clients.consumer.CooperativeStickyAssignor,\
  org.apache.kafka.clients.consumer.RangeAssignor
// Java 코드
Properties props = new Properties();
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group");
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
    CooperativeStickyAssignor.class.getName());

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("orders", "payments"));

Static Membership (KIP-345)

group.instance.id=consumer-1
  • consumer가 동일 ID로 재접속하면 partition 유지
  • 배포 / 재시작 시 rebalance 회피
  • 보통 수십 초 ~ 수 분 안의 짧은 재시작 가정
flowchart LR
    C1["Consumer 1\ngroup.instance.id=c1"] -->|"재시작"| C1b["Consumer 1 (재시작)"]
    C1b -->|"동일 ID로 재접속"| GC["Group Coordinator"]
    GC -->|"partition 유지 (rebalance 없음)"| P["p0, p1 유지"]

Offset 관리

flowchart LR
    C[Consumer] -->|fetch| K[(Kafka)]
    C -->|commit| Offset["__consumer_offsets\n(internal topic)"]
    Note1["정기 commit 또는 명시적"]
모드안전성
enable.auto.commit=true (기본)중복 또는 손실 가능
명시적 sync commit안전
명시적 async commit + 마지막 syncbalanced

Auto commit 함정

1. fetch [100, 101, 102]
2. 처리 중 ...
3. 5초 후 auto commit (103까지 commit)
4. consumer crash (102 처리 미완)
5. 재시작 → 103부터 fetch → 102 손실!

명시적 commit 패턴

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        process(record);
    }
    // 배치 처리 완료 후 명시적 commit
    consumer.commitSync();
}

CAUTION

Auto commit은 손실 또는 중복의 함정. exactly-once가 필요하면 명시적 commit + idempotent 처리.

Exactly-once 처리 패턴

Kafka 자체의 EOS (Exactly-Once Semantics)는 Kafka → Kafka 구간만 보장. 외부 DB까지 포함하려면 추가 패턴 필요.

flowchart LR
    K[(Kafka)] -->|"fetch"| C[Consumer]
    C -->|"처리"| DB[(Database)]
    C -->|"offset commit"| K
    DB -->|"idempotency key 체크"| DB
// Transactional producer (Kafka → Kafka exactly-once)
producer.initTransactions();

try {
    producer.beginTransaction();
    producer.send(new ProducerRecord<>("output-topic", key, value));
    // offset commit을 트랜잭션에 포함
    producer.sendOffsetsToTransaction(offsets, consumerGroupMetadata);
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}
// DB까지 exactly-once: idempotency key 패턴
@Transactional
public void processRecord(ConsumerRecord<String, Order> record) {
    String idempotencyKey = record.topic() + "-" + record.partition() + "-" + record.offset();

    if (processedOffsets.existsByKey(idempotencyKey)) {
        return;   // 이미 처리됨
    }

    orderService.process(record.value());
    processedOffsets.save(idempotencyKey);
    // DB 트랜잭션 커밋 후 offset commit
}

Consumer Lag

flowchart LR
    Last["partition의 last offset (LEO): 1000"] --> Lag
    Curr["consumer current offset: 800"] --> Lag
    Lag["Lag = 200\nprocessing 못 따라옴"]
도구의미
kafka-consumer-groups.sh --describe즉시 확인
Kafka Lag ExporterPrometheus 메트릭
Burrow (LinkedIn)더 세련된 lag 분석
# lag 확인
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --group my-consumer-group \
  --describe

# 출력 예시
GROUP           TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
my-group        orders    0          800             1000            200
my-group        orders    1          950             1000            50

Lag 기반 Auto-Scaling

flowchart LR
    Lag["Consumer Lag"] -->|KEDA| HPA["Kubernetes HPA"]
    HPA --> Scale["Consumer pod 수 증가"]

KEDA가 Kafka lag을 Kubernetes HPA 지표로 사용. spike 자동 흡수.

# KEDA ScaledObject
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: kafka-consumer-scaler
spec:
  scaleTargetRef:
    name: order-consumer
  minReplicaCount: 1
  maxReplicaCount: 10   # partition 수 이하로
  triggers:
    - type: kafka
      metadata:
        bootstrapServers: kafka:9092
        consumerGroup: my-consumer-group
        topic: orders
        lagThreshold: "100"   # lag 100 이상이면 scale out

Heartbeat / Session

설정의미기본
session.timeout.ms이 시간 heartbeat 없으면 죽었다고 판단45s
heartbeat.interval.msheartbeat 보내는 주기3s
max.poll.interval.mspoll 사이 최대 시간5min

WARNING

max.poll.interval.ms 초과 = consumer 추방. 처리 시간이 길면 명시적 늘림 또는 워커 별도 스레드.

Consumer Group 성능 튜닝

설정기본권장 (고처리량)
fetch.min.bytes11024 (배치 효율)
fetch.max.wait.ms500ms500ms
max.poll.records500100~500 (처리 시간에 따라)
max.partition.fetch.bytes1MB1MB
enable.auto.committruefalse (명시적 commit)
# 고처리량 consumer 설정
fetch.min.bytes=1024
fetch.max.wait.ms=500
max.poll.records=200
enable.auto.commit=false
auto.offset.reset=earliest

흔한 함정

WARNING

  1. Rebalancing 폭주 = consumer 자주 추가/제거 → 처리 정지 반복. cooperative + static membership.
  2. auto.offset.reset=latest + 다운 = 다운 사이 메시지 손실. 보통 earliest.
  3. 너무 큰 batch = max.poll.records 너무 큼 → 처리 시간 길어 추방.
  4. 상태 있는 consumer + rebalance = partition이 다른 consumer로 가면서 처리 진행 정보 손실. 상태는 외부 store에.
  5. partition 수 < consumer 수 = 초과 consumer는 idle. partition 수를 consumer 최대 수 이상으로 설정.

관련 위키

이 글의 용어 (5개)
[Distributed] Kafka: 분산 로그, partition, consumer groupdistributed-systems
정의 Apache Kafka = 분산 commit log. 고처리량 (수백만 msg/s), 영속, 수평 확장. event-driven 아키텍처 의 de facto. 핵심 개념: …
[Distributed] Message Broker 비교: Kafka / RabbitMQ / NATS / SQS / Redis Streamsdistributed-systems
정의 Message Broker = 생산자(Producer)와 소비자(Consumer) 사이에서 메시지를 중계하는 미들웨어. 비동기 통신, 부하 분산, 시스템 디커플링의 핵심. …
[Pattern] Idempotency Keys: 중복 요청 안전 처리distributed-systems
정의 Idempotency = 같은 요청을 N번 보내도 결과가 1번과 동일. 분산 시스템 / 결제 / API 의 안전망. [!IMPORTANT] 네트워크는 항상 timeout /…
[Pattern] Outbox Pattern: DB + 메시지의 원자성distributed-systems
정의 Outbox Pattern = DB 변경 + 메시지 발행 의 원자성 보장. 이중 쓰기 (dual write) 문제 의 표준 해결. 문제: Dual Write | 시나리오 |…
[Redis] Pub/Sub vs Streams: 휘발 신호 vs 영속 로그database-internals
정의 - Pub/Sub ( / ): 지금 듣고 있는 구독자 에게만 메시지가 전달되는 휘발성 신호. 영속 없음, ACK 없음. fan-out. - Streams ( / / / ):…

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기