Streams 리밸런스 프로토콜

Streams 리밸런스 프로토콜

카프카 4.2부터 Kafka Streams는 브로커 주도(broker-driven)의 전용 리밸런싱 시스템을 제공해요. 기존에는 태스크 할당을 클라이언트가 계산했지만, 새 프로토콜에서는 브로커가 그 역할을 맡아요. 이 페이지에서는 Streams Rebalance Protocol(KIP-1071)이 뭔지, 어떻게 켜고, 뭐가 지원되고 뭐가 안 되는지, 구성과 마이그레이션 방법까지 정리해드릴게요.

출처: 문서

본문

Streams Rebalance Protocol은 Kafka Streams 애플리케이션을 위해 특별히 설계된 브로커 주도 리밸런싱 시스템이에요. 일반 컨슈머의 리밸런스 조정을 클라이언트에서 브로커로 옮긴 KIP-848의 패턴을 따라, KIP-1071이 이 모델을 Kafka Streams 워크로드로 확장해요.

개요

그룹의 모든 멤버가 관여하는 리밸런스 이벤트 동안 클라이언트가 클라이언트 측에서 새 할당을 계산하는 대신, 할당은 브로커에서 지속적으로 계산돼요. 컨슈머 그룹을 사용하는 대신, 스트림즈 애플리케이션은 브로커에 스트림즈 그룹(streams group) 으로 등록하며, 브로커가 스트림즈 애플리케이션 인스턴스의 조정에 필요한 모든 메타데이터를 관리·노출해요.

이 접근 방식은 Kafka Streams 조정을 KIP-848에서 컨슈머를 위해 도입된 현대적 브로커 주도 리밸런스 모델과 일치시켜, 스트림즈 전용 의미론과 메타데이터 관리를 가진 전용 그룹 유형을 제공해요.

이 버전에서 지원되는 것

현재 릴리스에서 사용 가능한 기능:

  • 핵심 Streams Group Rebalance Protocol: group.protocol=streams 구성이 전용 스트림즈 리밸런스 프로토콜을 활성화해요. 이것은 스트림즈 그룹을 컨슈머 그룹과 분리하고, 스트림즈 전용 그룹 멤버십 수명주기와 브로커상의 메타데이터 관리를 제공해요.
  • Sticky Task Assignor: 리밸런스 중 태스크 이동을 최소화하는 기본 태스크 할당 전략이 포함돼요.
  • 인터랙티브 쿼리 지원: IQ 연산이 새 스트림즈 프로토콜과 호환돼요.
  • 새 Admin RPC: StreamsGroupDescribe RPC가 컨슈머 그룹 정보와 별개인 스트림즈 전용 메타데이터를 제공하며, Admin 인터페이스를 통해 접근할 수 있어요.
  • CLI 통합: bin/kafka-streams-groups.sh 스크립트로 스트림즈 그룹을 나열·설명·삭제할 수 있어요.
  • 오프라인 마이그레이션: 모든 멤버를 종료하고 session.timeout.ms가 만료되기를 기다린 후(또는 명시적 그룹 이탈을 강제한 후), 클래식 그룹을 스트림즈 그룹으로, 스트림즈 그룹을 클래식 그룹으로 변환할 수 있어요. 보존될 유일한 브로커 측 그룹 데이터는 커밋된 오프셋이에요. 내부 토픽(체인지로그와 리파티션 토픽)은 일반 카프카 토픽으로 계속 존재해요.

이 버전에서 지원되지 않는 것

아직 사용할 수 없어 새 프로토콜을 사용할 때 피해야 할 기능:

  • 정적 멤버십 (Static Membership): 클라이언트 instance.id를 설정하면 거부돼요.
  • 토폴로지 업데이트: 토폴로지가 크게 변경되면(예: 새 소스 토픽을 추가하거나 서브토폴로지 수 변경), 새 스트림즈 그룹을 만들어야 해요.
  • 고가용성 Assignor: sticky assignor만 지원돼요. 즉 "warmup tasks"와 랙 인지 할당은 아직 지원되지 않아요.
  • 정규식: 패턴 기반 토픽 구독은 지원되지 않아요.
  • 온라인 마이그레이션: 애플리케이션이 실행 중인 동안의 그룹 마이그레이션은 클래식과 새 스트림즈 프로토콜 사이에서 가능하지 않아요.
  • 커스텀 클라이언트 공급자: 커스텀 KafkaClientSupplier를 사용하면 restore/global 컨슈머, 프로듀서, admin 클라이언트만 제공할 수 있어요. "streams" 그룹이 활성화되면 "main" 컨슈머를 제공할 수 없어요.

왜 Streams Rebalance Protocol을 쓰나요?

Streams Rebalance Protocol은 클래식 클라이언트 주도 프로토콜보다 몇 가지 핵심 이점을 제공해요:

  • 브로커 주도 조정: 태스크 할당 로직을 클라이언트 대신 브로커에 중앙화해요. 단일 조정 지점에서 일관되고 권위 있는 태스크 할당 결정을 제공하고, 분할-뇌(split-brain) 시나리오의 가능성을 줄여요.
  • 더 빠르고 안정적인 리밸런스: 전역 동기화 지점을 제거해 리밸런스 기간과 영향을 줄여요. 멤버십 변경이나 장애 중 애플리케이션 다운타임을 최소화해요.
  • 더 나은 관찰성: 스트림즈 그룹과 컨슈머 그룹을 구분하는 전용 메트릭과 admin 인터페이스를 제공해, 브로커 측 관찰성으로 더 명확한 문제 해결을 가능하게 해요. 자세한 내용은 스트림즈 그룹 메트릭 문서를 참고해요.

프로토콜 활성화

Streams Rebalance Protocol은 Apache Kafka 4.2부터 새 클러스터에서 기본으로 활성화돼요. 이 프로토콜을 사용하려면 브로커와 클라이언트 모두 Apache Kafka 4.2 이상을 실행해야 해요.

브로커 구성

이 프로토콜은 새 Apache Kafka 4.2 클러스터에서 기본으로 활성화돼요. 기존 클러스터에서(4.2로 업그레이드한 후) 기능을 활성화하거나 명시적으로 제어하려면:

  • 기능 활성화:
bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature streams.version=1
  • 기능 비활성화:
bin/kafka-features.sh --bootstrap-server localhost:9092 downgrade --feature streams.version=0

클라이언트 구성

Kafka Streams 애플리케이션 구성에 다음을 설정해요:

group.protocol=streams

구성

브로커 구성

다음 브로커 구성이 스트림즈 그룹의 동작을 제어해요. 전체 세부사항은 브로커 구성 문서를 참고해요.

  • group.coordinator.rebalance.protocols: 활성화된 리밸런스 프로토콜의 목록. 스트림즈 그룹을 활성화하려면 "streams"가 프로토콜 목록에 포함되어 있어요.
  • group.streams.session.timeout.ms: 스트림즈 그룹 프로토콜을 사용할 때 클라이언트 장애를 감지하는 모든 스트림즈 그룹의 기본 타임아웃 (특정 스트림즈 그룹에 대해 명시적으로 오버라이드되지 않은 경우).
  • group.streams.min.session.timeout.ms: 최소 세션 타임아웃.
  • group.streams.max.session.timeout.ms: 최대 세션 타임아웃.
  • group.streams.heartbeat.interval.ms: 멤버에게 주어지는 기본 하트비트 간격.
  • group.streams.min.heartbeat.interval.ms: 최소 하트비트 간격.
  • group.streams.max.heartbeat.interval.ms: 최대 하트비트 간격.
  • group.streams.max.size: 단일 스트림즈 그룹이 수용할 수 있는 최대 스트림즈 클라이언트 수.
  • group.streams.num.standby.replicas: 각 태스크의 기본 스탠바이 복제본 수.
  • group.streams.max.standby.replicas: 스탠바이 복제본 구성의 동적 구성에 대한 최대값.
  • group.streams.initial.rebalance.delay.ms: 새(즉 이전에 비어 있던) 그룹의 첫 리밸런스는 더 많은 멤버가 그룹에 합류하도록 이 값만큼 지연돼요.

그룹 구성

리소스 유형 GROUP에 대한 구성은 DescribeConfigsIncrementalAlterConfigs에서 사용 가능하며, 특정 그룹에 대해 기본 브로커 구성을 동적으로 오버라이드해요. Admin Java 인터페이스나 bin/kafka-configs.sh 유틸리티로 설정할 수 있어요. 전체 세부사항은 그룹 구성 문서를 참고해요.

스트림즈 그룹에 사용 가능한 그룹 수준 구성:

  • streams.session.timeout.ms: 스트림즈 그룹 프로토콜을 사용할 때 클라이언트 장애를 감지하는 타임아웃.
  • streams.heartbeat.interval.ms: 멤버에게 주어지는 하트비트 간격.
  • streams.num.standby.replicas: 각 태스크의 스탠바이 복제본 수.
  • streams.initial.rebalance.delay.ms: 그룹의 첫 리밸런스는 더 많은 멤버가 그룹에 합류하도록 이 값만큼 지연돼요.

예시: 그룹 수준 구성 설정

bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type groups --entity-name wordcount \
  --add-config streams.num.standby.replicas=1

참고: 스트림즈 리밸런스 프로토콜에서 session.timeout.ms, heartbeat.interval.ms, num.standby.replicas는 그룹 수준 구성이며, 클라이언트 측에서 설정하면 무시돼요. 위와 같이 bin/kafka-configs.sh 도구를 사용해 이 구성들을 설정하세요.

Streams 구성

모든 Kafka Streams 구성에 대한 전체 세부사항은 스트림즈 구성 문서를 참고해요.

다음 클라이언트 구성이 스트림즈 리밸런스 프로토콜을 활성화해요:

  • group.protocol: 스트림즈 리밸런스 프로토콜을 사용할지 여부를 나타내는 플래그. 활성화하려면 streams로 설정 (기본값은 classic).

무시되는 구성

스트림즈 리밸런스 프로토콜이 활성화되면 다음 구성은 무시돼요:

  • acceptable.recovery.lag
  • max.warmup.replicas
  • num.standby.replicas (대신 그룹 수준 구성 사용)
  • probing.rebalance.interval.ms
  • rack.aware.assignment.tags
  • rack.aware.assignment.strategy
  • rack.aware.assignment.traffic_cost
  • rack.aware.assignment.non_overlap_cost
  • task.assignor.class
  • session.timeout.ms (대신 그룹 수준 구성 사용)
  • heartbeat.interval.ms (대신 그룹 수준 구성 사용)

관리

Admin API

Admin 인터페이스의 "streams groups" 메서드를 사용해 스트림즈 그룹을 프로그래매틱하게 관리해요. 이 API는 대부분 컨슈머 그룹 API와 같은 구현에 의해 지원돼요.

컨슈머 그룹 API와의 주요 차이점:

  • describeStreamsGroups는 DescribeStreamsGroup RPC를 사용하며 컨슈머 그룹과 다른 정보를 포함해요.
  • 스트림즈 그룹에는 추가 상태 NOT_READY가 있고, 클래식 프로토콜의 레거시 상태는 없어요.
  • removeMembersFromConsumerGroup은 이 버전에서 대응하는 API가 없어요. 클래식 컨슈머 그룹에 대해 LeaveGroup RPC를 사용하며, 이는 KIP-848 스타일 그룹에는 사용할 수 없기 때문이에요.

kafka-streams-groups.sh

스트림즈 그룹을 다루기 위한 새 도구 bin/kafka-streams-groups.sh가 추가됐어요. 이것은 스트림즈 그룹에 대해 bin/kafka-streams-application-reset.sh를 대체하며, 스트림즈 그룹을 나열·설명·삭제하는 데 사용할 수 있어요. 자세한 사용법은 kafka-streams-groups.sh 문서를 참고해요.

아키텍처와 동작 방식

Streams Groups

이 프로토콜은 컨슈머 그룹과 평행한 streams group의 개념을 도입해요. 스트림즈 클라이언트는 전용 하트비트 RPC인 StreamsGroupHeartbeat를 사용해 그룹에 합류하고, 그룹을 떠나며, 현재 소유한 태스크와 클라이언트별 메타데이터에 대해 그룹 코디네이터를 업데이트해요.

그룹 코디네이터는 컨슈머 그룹과 유사하게 스트림즈 그룹을 관리해요. 하트비트 응답을 통해 그룹 멤버 메타데이터를 지속적으로 업데이트하고, 변경이 감지되면 할당 로직을 실행해요. 그룹 코디네이터에 streams라는 새 그룹 유형이 도입되는데, 그룹 메타데이터, 토폴로지 메타데이터, 그룹 멤버 메타데이터에 대한 새 레코드 키와 값 유형이 있어요. 이 레코드는 __consumer_offsets 토픽에 영속화돼요.

그룹은 해당 GroupId를 사용하는 첫 하트비트 요청에 의해 정의되는 대로, streams 그룹, share 그룹, 또는 컨슈머 그룹일 수 있어요.

토폴로지 구성과 검증

스트림즈 클라이언트 간에 태스크를 할당하기 위해, 그룹 코디네이터는 멤버가 그룹에 합류할 때 초기화되고 컨슈머 오프셋 토픽에 영속화되는 토폴로지 메타데이터를 사용해요.

멤버가 스트림즈 그룹에 합류할 때마다 첫 하트비트 요청에 토폴로지 메타데이터가 포함돼요. 메타데이터는 토폴로지를 서브토폴로지의 집합으로 설명하는데, 각 서브토폴로지는 고유한 문자열 식별자로 식별되고 내부 토픽 생성과 할당에 관련된 메타데이터를 포함해요.

토폴로지 검증과 NOT_READY 상태

스트림즈 그룹 하트비트 처리 중에 그룹 코디네이터는 토폴로지가 필요로 하는 소스/싱크 또는 내부 토픽이 존재하지 않거나, 토폴로지가 성공적으로 실행되기 위해 필요한 구성과 다르다는 것을 감지할 수 있어요. 이것은 그룹 코디네이터가 다음 단계를 수행하는 "토폴로지 구성" 프로세스를 촉발해요:

  • 구성된 모든 소스 토픽이 존재하는지 확인.
  • "copartition groups"가 충족되는지 확인 — 즉 copartition되어야 하는 모든 소스 토픽이 실제로 copartition되었는지.
  • 소스 토픽 구성에서 모든 내부 토픽의 필요한 파티션 수를 도출.
  • 모든 내부 토픽이 올바른 구성으로 존재하는지 확인.

소스 토픽이나 내부 토픽이 없으면 그룹은 NOT_READY 상태로 들어가요. NOT_READY에서 모든 하트비트는 평소처럼 처리되며(보통 실패하지 않아야 함), 하트비트 응답의 상태는 어떤 종류의 문제가 존재하는지 나타내요. 그룹이 NOT_READY 상태일 때 모든 멤버는 빈 할당을 받아요.

중앙 집중식 할당 구성

핵심 할당 옵션은 각 클라이언트의 구성에 의존하지 않고 브로커에서 중앙으로 구성돼요. 이것은 스트림즈 애플리케이션을 재배포하지 않고 스트림즈 그룹을 튜닝할 수 있게 해줘요. 브로커 측에 도입된 핵심 할당 옵션은 num.standby.replicas예요. 이것은 브로커에서 전역적으로 구성할 수 있고, IncrementalAlterConfigsDescribeConfigs RPC를 통해 특정 스트림즈 그룹에 대해 동적으로 구성할 수 있어요.

마지막으로 사용된 할당 구성은 브로커의 그룹 메타데이터에 저장돼요. 이렇게 하면 할당 구성이 동적으로 변경되면 재할당이 즉시 촉발될 수 있어요.

모니터링과 메트릭

기존 그룹 메트릭은 스트림즈 그룹과 컨슈머 그룹을 구분하고 스트림즈 그룹 상태를 반영하도록 확장돼요. 전체 세부사항은 스트림즈 그룹 메트릭 문서를 참고해요.

프로토콜별 그룹 수

프로토콜 유형에 따른 그룹 수. 프로토콜 목록이 protocol=streams 변형으로 확장돼요:

kafka.server:type=group-coordinator-metrics,name=group-count,protocol={consumer|classic|streams}

상태별 Streams Group 수

상태에 따른 스트림즈 그룹 수:

kafka.server:type=group-coordinator-metrics,name=streams-group-count,state={empty|not_ready|assigning|reconciling|stable|dead}

Streams Group 리밸런스

스트림즈 그룹 리밸런스 센서:

kafka.server:type=group-coordinator-metrics,name=streams-group-rebalance-rate
kafka.server:type=group-coordinator-metrics,name=streams-group-rebalance-count

클래식 프로토콜에서의 마이그레이션

현재는 오프라인 마이그레이션만 지원돼요. Kafka Streams 애플리케이션을 클래식 프로토콜에서 스트림즈 리밸런스 프로토콜로 마이그레이션하려면:

  1. 모든 애플리케이션 인스턴스를 종료.
  2. session.timeout.ms가 만료되어 그룹이 비워질 때까지 기다림 (또는 명시적 그룹 이탈을 강제).
  3. 애플리케이션 구성을 group.protocol=streams로 업데이트.
  4. 애플리케이션 인스턴스를 재시작.

보존될 유일한 브로커 측 그룹 데이터는 커밋된 오프셋이에요. 다른 모든 그룹 메타데이터는 애플리케이션이 새 프로토콜로 시작할 때 재생성돼요. 내부 토픽(체인지로그와 리파티션 토픽)은 일반 카프카 토픽으로 계속 존재해요.

마찬가지로 같은 과정을 따르되 group.protocol=classic을 설정하면 스트림즈 그룹을 클래식 그룹으로 되돌릴 수 있어요.

경고: 온라인 마이그레이션(애플리케이션이 실행 중일 때 마이그레이션)은 이 버전에서 사용할 수 없어요. 프로토콜 사이를 마이그레이션할 때는 유지보수 창을 계획하세요.

경고: 오프라인 마이그레이션 코드의 심각한 브로커 측 버그(KAFKA-20254) 때문에, 4.2.0에서 클래식에서 스트림즈 그룹으로의 마이그레이션을 하지 않는 것이 좋아요. 새로 생성된 스트림즈 그룹은 영향을 받지 않아요. 수정 사항은 4.2.1에서 제공돼요.

더 알아보기