Kafka Streams Groups Tool

Kafka Streams Groups Tool

스트림즈 리밸런스 프로토콜(KIP-1071)을 쓰게 되면, 그룹을 관리할 때 쓰는 전용 CLI가 따로 있어요. 바로 kafka-streams-groups.sh예요. 그룹을 나열하고, 상태를 보고, 입력 토픽의 오프셋과 lag를 확인하고, 오프셋 리셋·삭제나 그룹 삭제 같은 운영 작업까지 할 수 있어요. 이 페이지에서 어떻게 쓰는지 정리해드릴게요.

출처: 문서

본문

kafka-streams-groups.sh를 사용해 Streams Rebalance Protocol(KIP‑1071)용 Streams 그룹을 관리해요: 그룹 나열·설명(describe), 멤버와 오프셋/lag 검사, 입력 토픽의 오프셋 리셋·삭제, 그룹 삭제(선택적으로 내부 토픽 포함) 등을 할 수 있어요.

개요

Streams 그룹은 Kafka Streams를 위한 브로커 조정(broker‑coordinated) 그룹 유형이에요. 클래식 컨슈머 그룹과는 구별되는 Streams 전용 RPC와 메타데이터를 사용해요. CLI는 Streams 전용 상태, 할당, 입력 토픽 오프셋을 표면화해 가시성과 관리를 단순화해요.

주의해서 사용하세요: 변경 작업(오프셋 리셋/삭제, 그룹 삭제)은 애플리케이션이 재시작할 때 데이터를 어떻게 재처리할지에 영향을 줘요. 항상 실행 전에 --dry-run으로 미리보기하고, 명령을 실행하기 전에 애플리케이션 인스턴스가 중지/비활성 상태이고 그룹이 비어 있는지 확인하세요.

Streams Groups 도구가 하는 일

  • 클러스터 전반의 Streams 그룹을 나열하고, 그룹 상태(Empty, Not Ready, Assigning, Reconciling, Stable, Dead)로 표시하거나 필터링해요.
  • Streams 그룹을 describe하고 표시해요: 그룹 상태, 그룹 epoch, 대상 할당 epoch(--state, 추가 세부정보는 --verbose).
  • 멤버별 정보: epoch, 현재 vs 대상 할당, 멤버가 여전히 클래식 프로토콜을 사용하는지 여부(--members--verbose).
  • 입력 토픽 오프셋과 lag(--offsets) — 처리가 얼마나 뒤처져 있는지 파악해요.
  • Streams 그룹의 입력 토픽 오프셋을 리셋해 정밀한 지정자(earliest, latest, to‑offset, to‑datetime, by‑duration, shift‑by, from‑file)로 재처리 경계를 제어해요. --dry-run 또는 --execute가 필요하며 비활성 인스턴스가 필요해요.
  • 입력 토픽의 오프셋을 삭제해 다음 시작 시 강제로 다시 소비하게 해요.
  • Streams 그룹을 삭제해 브로커 측 Streams 메타데이터(오프셋, 토폴로지, 할당)를 정리해요. 선택적으로 --internal-topics를 사용해 내부 토픽의 전체 또는 일부를 동시에 삭제할 수 있어요.

사용법

스크립트는 bin/kafka-streams-groups.sh에 있으며 --bootstrap-server를 통해 클러스터에 연결해요. 보안이 적용된 클러스터에서는 --command-config로 AdminClient 속성을 전달해요.

$ kafka-streams-groups.sh --bootstrap-server <host:port> [COMMAND] [OPTIONS]

참고: kafka-streams-groups.sh는 Streams 그룹용 Streams Admin API를 보완해요. CLI는 컨슈머 그룹 도구와 비슷한 정신으로 list/describe/delete 연산과 오프셋 관리를 노출하지만, KIP‑1071에 정의된 Streams 그룹에 맞게 조정되어 있어요.

명령어

Streams 그룹 나열

그룹 발견하기:

# 모든 Streams 그룹 나열
kafka-streams-groups.sh --bootstrap-server localhost:9092 --list

Streams 그룹 describe

그룹의 상태, 멤버, lag 검사하기:

# 그룹 describe: 상태 + epoch
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --describe --group my-streams-app --state --verbose

# 그룹 describe: 멤버 (할당 vs 대상, classic/streams)
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --describe --group my-streams-app --members --verbose

# 그룹 describe: 입력 토픽 오프셋과 lag
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --describe --group my-streams-app --offsets

입력 토픽 오프셋 리셋 (미리보기 → 적용)

모든 애플리케이션 인스턴스가 중지/비활성 상태인지 확인하세요. --execute를 사용하기 전에 항상 --dry-run으로 변경을 미리보세요.

# 모든 입력 토픽을 특정 타임스탬프로 리셋하는 것을 미리보기
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --group my-streams-app \
  --reset-offsets --all-input-topics --to-datetime 2025-01-31T23:57:00.000 \
  --dry-run

# 리셋 적용
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --group my-streams-app \
  --reset-offsets --all-input-topics --to-datetime 2025-01-31T23:57:00.000 \
  --execute

오프셋 삭제로 강제 재소비

모든 또는 특정 입력 토픽의 오프셋을 삭제해 그룹이 재시작 시 데이터를 다시 읽도록 해요.

# 모든 입력 토픽의 오프셋 삭제 (실행)
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --group my-streams-app \
  --delete-offsets --all-input-topics --execute

# 특정 토픽의 오프셋 삭제
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --group my-streams-app \
  --delete-offsets --topic input-a --topic input-b --execute

Streams 그룹 삭제 (정리)

그룹에 대한 브로커 측 Streams 메타데이터를 삭제하고 선택적으로 내부 토픽의 일부를 제거해요.

# Streams 그룹 메타데이터 삭제
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --delete --group my-streams-app

# 그룹과 함께 내부 토픽 일부 삭제 (주의해서 사용)
kafka-streams-groups.sh --bootstrap-server localhost:9092 \
  --delete --group my-streams-app \
  --internal-topics my-app-repartition-0,my-app-changelog

모든 옵션과 플래그

핵심 동작

  • --list: Streams 그룹 나열. --state로 상태에 따라 표시/필터링.
  • --describe: --group으로 선택한 그룹을 describe. 결합 가능: --state(그룹 상태와 epoch), --members(멤버와 할당), --offsets(입력 및 리파티션 토픽 오프셋/lag). 추가 세부정보(예: 해당되는 리더 epoch)는 --verbose.
  • --reset-offsets: 입력 토픽 오프셋 리셋 (한 번에 한 그룹; 인스턴스는 비활성이어야 함). 정확히 하나의 지정자 선택: --to-earliest, --to-latest, --to-current, --to-offset <n>, --by-duration <PnDTnHnMnS>, --to-datetime <YYYY-MM-DDTHH:mm:SS.sss>, --shift-by <n>(±), --from-file(CSV)
    • 범위: --all-input-topics 또는 하나/여러 개의 --topic <name>; 일부 빌드는 --all-topics(브로커 토폴로지 메타데이터 기준 모든 입력 토픽)도 지원해요.
    • 안전: --dry-run 또는 --execute가 필요해요.
  • --delete-offsets: --all-input-topics, 특정 --topic 이름, 또는 --from-file에 대한 오프셋 삭제.
  • --delete: Streams 그룹 메타데이터 삭제; 선택적으로 --internal-topics <list>를 전달해 내부 토픽의 일부를 삭제.

공통 플래그

  • --group <id>: 대상 Streams 그룹(application.id).
  • --all-groups: 모든 그룹에 대해 동작 (--delete와 함께 허용).
  • --bootstrap-server <host:port>: 연결할 브로커 (필수).
  • --command-config <file>: AdminClient 속성(보안, 타임아웃 등).
  • --timeout <ms>: 일부 연산에서 그룹 안정화를 기다리는 시간 (기본값: 5000ms).
  • --dry-run, --execute: 변경 작업의 미리보기 vs 적용.
  • --help, --version, --verbose: 사용법, 버전, 상세 출력.

모범 사례와 안전

  • --execute 전에 --dry-run으로 변경을 미리보아 토픽 범위와 영향을 확인하세요.
  • --internal-topics는 주의해서 사용하세요: 내부 토픽을 삭제하면 상태를 뒷받침하는 토픽이 제거돼요. 입력 토픽에서 상태를 재구축할 의도가 있을 때만 하세요.

이 페이지는 KIP‑1071에 정의되고 Apache Kafka에 구현된 Streams 그룹에 대한 kafka-streams-groups.sh의 기능을 문서화해요.

더 알아보기