업그레이드 가이드

업그레이드 가이드

카프카 스트림즈를 더 새 버전으로 올릴 때는 몇 가지 주의할 점이 있어요. 특히 여러 번 바운스(롤링 업그레이드)해야 하는 경우, 다운그레이드 시 주의사항, 그리고 버전마다 달라진 API 변경점을 알아야 해요. 이 페이지에서 업그레이드 절차와 과거 릴리스들의 주요 호환성·API 변경을 정리해드릴게요.

출처: 문서

본문

업그레이드 가이드와 API 변경

어떤 이전 버전에서든 4.3.0으로의 업그레이드가 가능해요. 3.4 이하에서 업그레이드한다면 두 번의 롤링 바운스(rolling bounce)가 필요하며, 첫 번째 롤링 바운스 단계에서 upgrade.from="older version" 구성(가능한 값은 "2.4" - "3.4")을 설정하고 두 번째 단계에서는 제거해야 해요. 이는 두 가지 변경을 안전하게 처리하기 위해 필요해요. 첫 번째는 외래 키 조인 직렬화 형식의 변경이고 두 번째는 내부 리파티션 토픽의 직렬화 형식 변경이에요. 자세한 내용은 KIP-904를 참고해요:

  1. 롤링 바운스를 위해 애플리케이션 인스턴스를 준비하고 구성 upgrade.from이 업그레이드하는 버전으로 설정되어 있는지 확인.
  2. 애플리케이션의 각 인스턴스를 한 번 바운스.
  3. 새로 배포된 4.3.0 애플리케이션 인스턴스를 두 번째 롤링 바운스 라운드를 위해 준비하고 구성 upgrade.from의 값을 제거.
  4. 업그레이드를 완료하기 위해 애플리케이션의 각 인스턴스를 한 번 더 바운스.

대안으로 오프라인 업그레이드도 가능해요. 0.11.0.x만큼 오래된 어떤 버전에서도 4.3.0으로 오프라인 모드로 업그레이드하려면 다음 단계가 필요해요:

  1. 모든 이전(예: 0.11.0.x) 애플리케이션 인스턴스를 중지.
  2. 코드를 업데이트하고 이전 코드·jar 파일을 새 코드·새 jar 파일로 교체.
  3. 모든 새(4.3.0) 애플리케이션 인스턴스를 재시작.

카프카 브로커 버전과 Streams API 호환성을 보여주는 표는 Broker Compatibility를 참고해요.

지난 릴리스들의 주목할 만한 호환성 변경

  • 4.0.0부터 Kafka Streams는 2.1 이상 브로커에 대해 실행할 때만 호환돼요. 또한 정확히 한 번 의미론(EOS)은 브로커가 적어도 2.5 버전이기를 요구해요.
  • 3.5.x 이상에서 3.4.x 이하로의 다운그레이드는 특별한 주의가 필요해요: 3.5.0 릴리스부터 Kafka Streams는 리파티션 토픽에 새 직렬화 형식을 사용해요. 즉 이전 버전의 Kafka Streams는 더 새 버전이 쓴 바이트를 인식하지 못하므로, 실행 중에 3.5.0 이상의 Kafka Streams를 이전 버전으로 다운그레이드하기가 더 어려워요. KIP-904를 참고해요. 다운그레이드하려면 먼저 구성 "upgrade.from"을 다운그레이드 대상 버전으로 전환해요. 이것은 애플리케이션에서 새 직렬화 형식 쓰기를 비활성화해요. 새 직렬화 형식으로 리파티션 토픽에 쓰인 "in-flight" 메시지의 처리를 애플리케이션이 끝냈는지 확인하기 위해 이 상태에서 충분히 기다리는 것이 중요해요. 그 후 애플리케이션을 3.5.x 이전 버전으로 다운그레이드할 수 있어요.
  • 3.0.x 이상에서 2.8.x 이하로의 다운그레이드는 특별한 주의가 필요해요: 3.0.0 릴리스부터 Kafka Streams는 온디스크 형식이 변경된 더 새 RocksDB 버전을 사용해요. 즉 이전 버전 RocksDB는 그 더 새 버전 RocksDB가 쓴 바이트를 인식하지 못하므로, 실행 중에 3.0.0 이상의 Kafka Streams를 이전 버전으로 다운그레이드하기가 더 어려워요. 사용자는 이전 버전의 Kafka Streams 바이트코드로 교체하기 전에 먼저 새 버전 Kafka Streams가 쓴 로컬 RocksDB 상태 저장소를 삭제해야 하며, 그런 다음 이전 온디스크 형식으로 상태 저장소를 체인지로그에서 복원할 거예요.
  • 4.0.x 이상에서 4.0.0 이전 버전으로의 다운그레이드는 특별한 주의가 필요해요: 4.0.0 릴리스부터 Kafka Streams는 RocksDB를 7.9.2에서 9.7.3으로 업그레이드했어요. 이 업그레이드는 RocksDB 파일 형식 버전을 5에서 6(RocksDB 8.6에서 도입)으로 변경해요. 더 새 RocksDB 버전(9.7.3)은 이전 형식(버전 5)으로 쓰인 상태 저장소를 읽을 수 있지만, 더 이전 RocksDB 버전은 새 형식(버전 6)으로 쓰인 상태 저장소를 읽을 수 없어요. 즉 4.0.x 이상 Kafka Streams에서 4.0.0 이전 버전으로의 실행 중 다운그레이드는 기본으로 가능하지 않아요. 사용자는 4.0.0 이전 버전으로 다운그레이드하기 전에 먼저 4.0.x 이상 Kafka Streams가 쓴 로컬 RocksDB 상태 저장소를 삭제해야 하며, 그런 다음 이전 파일 형식으로 상태 저장소를 체인지로그에서 복원할 거예요.
  • 4.3.x 이상에서 4.2.x 이하로의 다운그레이드는 특별한 주의가 필요해요: 4.3.0 릴리스부터 Kafka Streams는 상태 저장소 체인지로그 오프셋을 태스크별 .checkpoint 파일이 아닌 각 상태 저장소 내부에 영속해요 (KIP-1035). 내장 RocksDB 저장소의 경우 오프셋이 각 RocksDB 인스턴스 내부의 전용 offsets 칼럼 패밀리에 쓰여요. 이전 Kafka Streams 버전은 RocksDB를 열 때 이 칼럼 패밀리를 선언하지 않으므로 저장소를 열지 못하고 애플리케이션이 시작 시 충돌할 거예요. 따라서 실행 중 다운그레이드는 지원되지 않아요. 4.2.x 이하로 다운그레이드하려면 각 애플리케이션 인스턴스를 중지하고 로컬 상태 디렉터리(state.dir)를 삭제한 다음 이전 버전을 시작해요 — Kafka Streams가 이전 온디스크 형식으로 모든 상태 저장소를 체인지로그 토픽에서 복원할 거예요.
  • Kafka Streams는 같은 물리적 상태 디렉터리에서 다른 프로세스로 같은 애플리케이션의 여러 인스턴스를 실행하는 것을 지원하지 않아요. 2.8.0(뿐만 아니라 2.7.1, 2.6.2)부터 이 제한이 시행돼요. Kafka Streams의 한 개 이상의 인스턴스를 실행하려면 state.dir에 다른 값으로 구성해야 해요.
  • 2.6.x부터 새 처리 모드인 EOS version 2를 사용할 수 있어요. 이것은 3.0+ 애플리케이션 버전의 경우 "processing.guarantee""exactly_once_v2"로 설정하거나, 2.62.8 버전의 경우 "exactly_once_beta"로 설정해 구성할 수 있어요. 이 새 기능을 사용하려면 브로커가 2.5.x 이상이어야 해요. 이전 버전에서 EOS 애플리케이션을 업그레이드하고 3.0+에서 이 기능을 활성화하려면, 먼저 애플리케이션을 3.0.x로 업그레이드하면서 "exactly_once"를 유지하고, 그런 다음 두 번째 롤링 바운스 라운드를 해 "exactly_once_v2"로 전환해야 해요. 2.62.8 버전 사이에서 "exactly_once_beta" 구성으로 업그레이드하려면 같은 단계를 따르되 구성을 "exactly_once_beta"로 해요. 2.6+에서 3.0 이상으로 "exactly_once_beta"를 사용하는 애플리케이션을 업그레이드하는 데는 특별한 단계가 필요 없어요: 롤링 업그레이드 중에 구성을 "exactly_once_beta"에서 "exactly_once_v2"로 바꾸면 돼요. 다운그레이드의 경우 반대로: 먼저 "exactly_once_v2"에서 "exactly_once"로 구성을 전환해 2.6.x 애플리케이션에서 기능을 비활성화해요. 그 후 애플리케이션을 2.6.x 이전 버전으로 다운그레이드할 수 있어요.
  • 2.6.0부터 Kafka Streams는 MacOS 10.14 이상을 요구하는 RocksDB 버전에 의존해요.

4.3.0의 Streams API 변경

참고: Kafka Streams 4.3.0에는 RocksDB 상태 저장소 레이어의 치명적인 네이티브 메모리 누수(KAFKA-20616)가 있어요. ColumnFamilyOptions(offsets 칼럼 패밀리용)가 닫히지 않고, 칼럼 패밀리 핸들이 close-path 예외에서 누출될 수 있어, 계단식 태스크 클로즈(예: 리밸런스나 오류 트리거 복구)에서 무한 오프힙 메모리 증가와 결국 OOM으로 이어져요. Kafka Streams를 실행하는 사용자는 이에 대한 수정을 포함하는 4.3.1로 직접 업그레이드하는 것을 고려해야 해요.

  • Kafka Streams는 이제 KIP-1270을 통해 글로벌 저장소/KTable 처리를 위한 ProcessingExceptionHandler를 지원해요. 이전에는 ProcessingExceptionHandler가 일반 스트림 태스크에만 적용됐어요. 이 릴리스에서 새 구성 processing.exception.handler.global.enabledtrue(권장)로 설정해 글로벌 저장소/KTable에 대한 예외 처리를 구성할 수 있어요. 활성화되면 구성된 ProcessingExceptionHandler가 글로벌 저장소/KTable 처리 중 발생하는 예외에 대해 호출돼요. DLQ(Dead Letter Queue) 지원은 글로벌 저장소/KTable에 아직 없고 향후 릴리스에서 추가될 예정이에요. 자세한 내용은 KIP-1270에서 찾을 수 있어요.
  • 스트림즈 스레드 메트릭 commit-ratio, process-ratio, punctuate-ratio, poll-ratio와 스트림 상태 업데이터 메트릭 active-restore-ratio, standby-restore-ratio, idle-ratio, checkpoint-ratio가 업데이트됐어요. 각 메트릭은 이제 롤링 측정 윈도우에 걸쳐 이 스레드가 주어진 액션({action})을 수행하는 시간의 비율을 그 윈도우의 총 경과 시간에 대해 보고해요. 유효 윈도우 기간은 메트릭 구성으로 결정돼요: metrics.sample.window.ms(샘플별 윈도우 길이)와 metrics.num.samples(롤링 윈도우 수).
  • Kafka Streams는 이제 애플리케이션 시작 중에 일정 기간 수정되지 않은 로컬 상태 디렉터리와 체크포인트 파일을 정리할 수 있게 해줘요. 새 state.cleanup.dir.max.age.ms 구성으로 구성할 수 있어요. KIP-1259 참조.
  • Kafka Streams는 이제 상태 저장소 체인지로그 오프셋을 단일 태스크별 .checkpoint 파일이 아닌 각 상태 저장소 내부에 영속해요 (KIP-1035). 이것은 내부 인프라 변경이며 대부분 사용자에게 투명해요 — 기존 태스크별 .checkpoint 파일은 첫 시작 시 자동으로 마이그레이션되고 애플리케이션이나 운영자 조치가 필요 없어요. EOS 크래시 동작은 4.3에서 변하지 않았어요: 상태 저장소는 여전히 지워지고 체인지로그에서 완전히 복원돼요. KIP-1035는 상태 저장소에 대한 트랜잭션 의미론을 만들어 EOS 상태 쓰기를 트랜잭션화하고 전체 복원을 건너뛸 KIP-892의 전제 조건이에요. 커스텀 StateStore 구현 작성자는 managesOffsets(), commit(Map<TopicPartition, Long>), committedOffset(TopicPartition)을 통해 자신의 오프셋 관리를 옵트인할 수 있어요. 다운그레이드 영향은 지난 릴리스들의 주목할 만한 호환성 변경을 참고해요.
  • KIP-1035의 일부로, 저장소별 체인지로그 오프셋은 각 커밋에 RocksDB에 기록되고, RocksDB가 메모테이블을 SST 파일로 플러시할 때(즉 메모테이블이 write_buffer_size(기본 16 MB)를 채우는 유기적 플러시 또는 정상 저장소 종료 시)에만 디스크에서 내구성이 생겨요. 이전 릴리스는 매 커밋마다 RocksDB를 강제 플러시했어요. 결과적으로 메모테이블이 거의 채워지지 않는 저트래픽 저장소의 경우 온디스크 오프셋이 다음 정상 종료까지 저장소의 실제 위치보다 뒤처질 수 있어요. 내구성 모델과 저트래픽 저장소의 플러시 빈도 튜닝 지침은 메모리 관리: RocksDB를 참고해요. 프로세스가 비정상적으로 종료되고(예: SIGKILL/OOM-킬, 또는 셧다운 유예 기간 내에 완료되지 않는 KafkaStreams#close) 체인지로그 리텐션·컴팩션이 그 이후 체인지로그의 로그 시작 오프셋을 그 오래된 오프셋을 넘어 진행시켰다면, 복원 컨슈머가 재시작 시 범위 밖을 찾아요 — OffsetOutOfRangeException/TaskCorruptedException으로 기록 — 태스크가 체인지로그에서 자동으로 재초기화돼요 (데이터 손실 없지만 전체 재-복원).

Processor API용 헤더 인지 상태 저장소 (KIP-1271)

Kafka Streams가 헤더 인지 상태 저장소를 추가해요. 이름이 WithHeaders로 끝나는 새 Stores 공급자와 일치하는 StoreBuilder 팩토리로 옵트인해요. 예:

  • persistentTimestampedKeyValueStoreWithHeaders + timestampedKeyValueStoreWithHeadersBuilder
  • persistentTimestampedWindowStoreWithHeaders + timestampedWindowStoreWithHeadersBuilder
  • persistentSessionStoreWithHeaders + sessionStoreWithHeadersBuilder

Processor API 상태 저장소 문서를 참고해요. 같은 헤더 없는 Stores 공급자·빌더를 계속 사용하는 기존 애플리케이션은 영향받지 않아요: 저장 형식, 체인지로그, 성능이 이전과 같아요. 헤더 인지 형식을 채택하는 저장소에 대해 KIP-1271은 단일 롤링 바운스 업그레이드를 정의해요: 체인지로그 토픽 형식은 변하지 않고, 레거시 행은 다시 쓰일 때까지 빈 헤더 집합으로 읽히며, RocksDB 기반 저장소는 접근 시 데이터를 지연 마이그레이션해요. 마이그레이션 후 제자리 다운그레이드는 로컬 저장소 데이터를 지우고 체인지로그에서 복원하지 않는 한 지원되지 않아요. 헤더 저장은 헤더 없는 저장소보다 디스크와 직렬화 비용을 증가시켜요.

TopologyTestDriver와 인터랙티브 쿼리가 새 저장소 유형을 지원해요. 기존 store() 파사드는 레코드 헤더를 노출하지 않고 값(또는 ValueAndTimestamp)을 계속 반환해요.

DSL 연산자용 헤더 인지 상태 저장소 (KIP-1285)

KIP-1285는 DSL 연산자가 KIP-1271이 도입한 헤더 인지 상태 저장소를 사용할 수 있게 해줘요. dsl.store.format=HEADERS로 설정해 지원되는 DSL 연산자에 헤더 인지 저장소를 사용해요. 이 저장소는 레코드 헤더를 값·타임스탬프와 함께 유지할 수 있어요.

이 구성은 상태 저장소 형식만 선택해요. DSL 연산자가 출력 레코드용 헤더를 어떻게 만드는지는 정의하지 않아요 (아래 현재 제한사항 참고).

// Enable headers-aware stores globally for all DSL operators
Properties props = new Properties();
props.put(StreamsConfig.DSL_STORE_FORMAT_CONFIG, "HEADERS");

연산자별 커스터마이즈는 Materialized.withStoreType(...)로 커스텀 DslStoreSuppliers를 제공하거나 명시적 헤더 인지 저장소 공급자를 공급해 가능해요. 기존 boolean isTimestamped 생성자와 DslKeyValueParams·DslWindowParamsisTimestamped() 메서드, 그리고 DslSessionParams의 3인자 생성자는 DslStoreFormat 기반 생성자를 위해 폐기됐어요. 기존 애플리케이션은 기본적으로 영향받지 않아요.

현재 제한사항

오늘날 DSL 결과 헤더는 다음과 같이 동작해요:

  • 집계(count, reduce, aggregate, 윈도우·세션 윈도우 변형 포함), KTable-KTable 조인(inner/left/outer), 구체화된 KTable.mapValues, KStream.toTable(), StreamsBuilder.table()은 구체화된 저장소에 빈 헤더를 써요.
  • KStream-KStream 조인 윈도우 저장소는 소스 레코드 헤더를 유지하지만, 조인 결과 레코드는 계산되거나 병합된 헤더를 가지지 않아요. 결과를 트리거한 레코드의 헤더를 가질 수 있어요.
  • suppress()와 left/outer 스트림-스트림 조인은 헤더 인지가 아닌 버퍼 저장소를 사용해요. 그 버퍼를 통과하는 레코드는 헤더를 잃어요.

후속 KIP가 사용자에게 DSL 결과 헤더가 어떻게 계산되는지 명시적으로 제어하는 방법을 제공할 거예요. 자세한 내용은 상태 저장 변환을 참고해요.

체인지로그, 마이그레이션, 성능

KIP-1285는 체인지로그 와이어 형식, 마이그레이션 절차, 저장소별 오버헤드를 변경하지 않아요 — 그것들은 기본 KIP-1271 저장소의 속성이에요. 체인지로그 호환성, DEFAULTHEADERS의 지연 키별 RocksDB 마이그레이션, 복원 동작, 레코드별 크기 영향에 대한 전체 설명은 위 KIP-1271 섹션을 참고해요. DSL 구성 dsl.store.format은 어떤 연산자가 참여하는지만 제어해요. 연산자가 헤더 인지 저장소를 사용하게 되면 저장소 런타임 동작은 Processor API 경우와 동일해요.

streams-scala 모듈 폐기 (KIP-1244)

kafka-streams-scala 모듈(org.apache.kafka.streams.scala 패키지)이 4.3.0에서 폐기되고 5.0에서 제거될 예정이에요. 코드 예시가 있는 상세 마이그레이션 가이드는 Streams Scala에서 Java API로 마이그레이션하기를 참고해요.

4.2.0의 Streams API 변경

Streams Rebalance Protocol 핵심 기능 세트의 GA (KIP-1071)

Streams Rebalance Protocol은 Kafka Streams 애플리케이션을 위해 특별히 설계된 브로커 주도 리밸런싱 시스템이에요. 이 릴리스는 KIP-1071에 상세히 설명된 핵심 기능의 General Availability를 표시해요. 기능 세트, 설계, 사용법, 마이그레이션에 대한 자세한 내용은 개발자 가이드를 참고해요.

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

기타 변경

  • Kafka Streams가 이제 DLQ(Dead Letter Queue)를 지원해요. 새 구성 errors.dead.letter.queue.topic.name(KIP-1034)으로 DLQ 토픽 이름을 지정할 수 있어요. 설정되고 DefaultProductionExceptionHandler를 사용하면 예외를 일으키는 레코드가 DLQ 토픽으로 전달돼요. 커스텀 예외 핸들러를 사용하면 DLQ 레코드를 만들어 보내는 것은 커스텀 핸들러의 몫이므로 구현에 따라 errors.dead.letter.queue.topic.name 구성을 무시할 수 있어요.
  • org.apache.kafka.streams.CloseOptions 클래스가 기존 KafkaStreams$CloseOptions를 대체해요 (KIP-1153). 후자는 폐기되고 다음 메이저 릴리스에서 제거돼요.
  • 새 구성 allow.os.group.write.access(KIP-1230)로 OS 그룹에 대한 쓰기 접근을 상태 저장소 디렉터리에 허용할 수 있어요.
  • org.apache.kafka.streams.errors.BrokerNotFoundException이 폐기되고 다음 메이저 릴리스에서 제거될 예정이에요 (KIP-1195).
  • 스트림즈 리밸런스 콜백 지연을 모니터링할 수 있는 리밸런스 리스너 메트릭(KIP-1216)이 추가됐어요.
  • Kafka Streams 클라이언트 상태 메트릭(client-state)에 application-id 태그가 추가됐어요 (KIP-1221).

4.1.0의 Streams API 변경

참고: Kafka Streams 4.1.0에는 범위 스캔과 특정 DSL 연산자(세션 윈도우, sliding 윈도우, 스트림-스트림 조인, 외래 키 조인) 사용자에게 영향을 주는 치명적 메모리 누수 버그(KAFKA-19748)가 있어요. 사용자는 이에 대한 수정을 포함하는 4.1.1로 직접 업그레이드하는 것을 고려해야 해요.

Streams Rebalance Protocol 조기 접근

Streams Rebalance Protocol은 Kafka Streams 애플리케이션을 위해 특별히 설계된 브로커 주도 리밸런싱 시스템이에요. 일반 컨슈머의 리밸런스 조정을 클라이언트에서 브로커로 옮긴 KIP-848의 패턴을 따라 KIP-1071이 이 모델을 Kafka Streams 워크로드로 확장해요. 그룹의 모든 멤버가 관여하는 리밸런스 이벤트 동안 클라이언트가 새 할당을 계산하는 대신 할당은 브로커에서 지속적으로 계산돼요. 컨슈머 그룹을 사용하는 대신 스트림즈 애플리케이션은 스트림즈 그룹으로 브로커에 등록하며, 브로커가 스트림즈 애플리케이션 인스턴스 조정에 필요한 모든 메타데이터를 관리·노출해요.

이 Early Access 릴리스는 KIP-1071에 상세힌 기능의 부분 집합을 다뤄요. 프로덕션에서 새 프로토콜을 사용하지 마세요. API는 향후 릴리스에서 변경될 수 있어요.

Early Access에 포함된 것: 핵심 Streams Group Rebalance Protocol(group.protocol=streams), Sticky Task Assignor, 인터랙티브 쿼리 지원, 새 Admin RPC(StreamsGroupDescribe), CLI 통합(kafka-streams-groups.sh).

Early Access에 포함되지 않은 것: 정적 멤버십, 토폴로지 업데이트, High Availability Assignor, 정규식, 리셋 연산, 프로토콜 마이그레이션.

기타 변경

  • KIP-1111: 토폴로지의 모든 내부 리소스(내부 토픽과 상태 저장소 포함)에 대해 명시적 명명을 강제할 수 있게 해줘요. StreamsConfig#ENSURE_EXPLICIT_INTERNAL_RESOURCE_NAMING_CONFIG로 활성화하며, true로 설정하면 내부 리소스에 자동 생성 이름이 있으면 애플리케이션이 시작을 거부해요.

4.0.0의 Streams API 변경

이 릴리스에서 eos-v1(Exactly Once Semantics version 1)은 더 이상 지원되지 않아요. eos-v2를 사용하려면 브로커가 2.5 이상이어야 해요. 또한 AK 3.5 릴리스까지의 모든 폐기된 메서드·클래스·API·구성 파라미터가 제거됐어요. 주요 사항:

  • Java·Scala 모두의 이전 processor API, KStream#through().
  • Java·Scala 모두의 "transformer" 메서드·클래스. (KStreams#transformValues()에서 KStreams.processValues()로의 마이그레이션은 KAFKA-19668 때문에 안전하지 않을 수 있어요. 마이그레이션 가이드 참조.)
  • kstream.KStream#branch, Time/Session/Join/SlidingWindows의 빌더 메서드, KafkaStreams#setUncaughtExceptionHandler().

추가 변경:

  • ClientInstanceIds 인스턴스가 KIP-714 ID용 글로벌 컨슈머 Uuid를 이전에는 글로벌 스트림-스레드 이름만이었던 키(이제 "-global-consumer"가 붙음)로 저장.
  • 구성 default.deserialization.exception.handlerdefault.production.exception.handler 폐기 (KIP-1056). 새 구성 deserialization.exception.handler, production.exception.handler 사용.
  • 이전 Processor API가 증분적으로 대체·폐기됨. KIP-1070이 MockProcessorContext, Transformer, TransformerSupplier, ValueTransformer, ValueTransformerSupplier를 폐기.
  • KIP-1065: 이제 ProductionExceptionHandler가 (재시도 가능한) TimeoutException에서 호출되며, 기본 핸들러는 기존 동작을 유지하기 위해 RETRY를 반환해요. 그러나 커스텀 핸들러는 이제 CONTINUEFAIL을 반환해 무한 재시도 루프를 끊을 수 있어요.
  • KIP-1076: 브로커 측에서 KIP-714 브로커 플러그인을 통해 Kafka Streams 메트릭을 수집할 수 있어요.
  • RocksDB 의존성을 7.9.2에서 9.7.3으로 업그레이드. API 변경(AccessHint 클래스 제거, BlockBasedTableConfig의 압축 블록 캐시 메서드 제거 등)은 주로 커스텀 RocksDB 구성을 구현하는 사용자에게 관련돼요. rocksdb.config.setter를 사용하는 고급 RocksDB 커스터마이즈 사용자는 상세 RocksDB 9.7.3 체인지로그를 확인해 적응해야 해요.

3.9.0의 Streams API 변경

  • KIP-1033: 기록 처리 중 예외를 스트림즈 애플리케이션 밖으로 던지는 대신 처리 예외 핸들러를 제공할 수 있게 해줘요. StreamsConfig#PROCESSING_EXCEPTION_HANDLER_CLASS_CONFIG로 구성하며, 지정된 핸들러는 org.apache.kafka.streams.errors.ProcessingExceptionHandler 인터페이스를 구현해야 해요.
  • 새 구성 log.summary.interval.ms(KIP-1049)로 스트림-스레드 런타임 요약의 로깅 간격을 커스터마이즈할 수 있어요. 기본적으로 요약은 2분마다 기록돼요.

3.8.0의 Streams API 변경

  • task.assignor.class 구성(KIP-924)으로 커스터마이즈 가능한 태스크 할당 전략을 지원해요. 커스텀 태스크 assignor를 지정하지 않으면 기본 HighAvailabilityTaskAssignor가 사용돼요. internal.task.assignor.class 구성 사용자는 새 task.assignor.class로 전환해야 해요.
  • Processor API가 KIP-813로 추가된 소위 읽기 전용 상태 저장소를 지원해요. 이 저장소는 전용 체인지로그 토픽이 없고 KTable처럼 소스 토픽을 장애 허용에 사용해요.
  • 누출된 상태 저장소 반복자 감지 개선을 위해 새 메트릭 num-open-iterators, iterator-duration-avg, iterator-duration-max, oldest-iterator-open-since-ms를 추가 (KIP-989).

3.7.0의 Streams API 변경

  • KIP-988: KafkaStreams#setStandbyUpdateListener() 추가. 스탠바이 태스크 변경을 지속적으로 모니터링하는 StandbyUpdateListener 구현을 제공할 수 있어요.
  • IQv2 개선: RangeQuery가 내림차순/오름차순 데이터를 지원하고 (KIP-985), TimestampedKeyQuery·TimestampedRangeQuery(KIP-992), 버전 상태 저장소의 MultiVersionedKeyQuery(KIP-968)·VersionedKeyQuery(KIP-960) 지원 추가.
  • KIP-962: 카프카 스트림즈 조인 연산자의 non-null 키 요구사항 완화. left/outer 조인에서 null-키 레코드를 더 이상 버리지 않고 ValueJoinernull 값으로 호출. 기존 동작을 유지하려면 .filter((key, value) -> key != null) 연산자를 앞에 추가.
  • default.dsl.store 구성 폐기, 새 dsl.store.suppliers.class 구성 사용.
  • rack.aware.assignment.strategy에 새 구성 옵션 balance_subtopology 도입.

3.6.0의 Streams API 변경

  • KIP-925: 랙 인지 태스크 할당 도입. StickyTaskAssignor 또는 HighAvailabilityTaskAssignor에 대해 활성화해 특정 조건에서 교차 랙 트래픽을 최소화하는 태스크 할당을 계산할 수 있어요.
  • KIP-941: IQv2의 RangeQuerynull 상한·하한을 허용("경계 없음" 의미론).
  • KIP-923: KStream-KTable 조인에 유예 기간 추가 옵션.

3.5.0의 Streams API 변경

  • KIP-889·KIP-914: 버전 키-값 저장소(versioned key-value stores) 도입. 키당 단일 레코드 버전 대신 키당 여러 기록 버전을 저장할 수 있어, 지정된 타임스탬프 시점의 최신 레코드를 반환하는 timestamped 검색 연산을 지원해요. 옵트인만 가능해요. 글로벌 KTable은 지원되지 않고 suppress()와 함께 동작하지 않아요.
  • KIP-914: 버전 키-값 저장소를 사용하면 DSL 처리가 순서가 뒤집힌 데이터를 더 잘 처리할 수 있어요.
  • KIP-904: KTable 집계 구현 개선. 가짜 중간 결과를 피해 단일 업데이트만 적용.
  • KIP-399: ProductionExceptionHandler가 이제 직렬화 오류도 다뤄요.
  • KIP-907: 새 Serde 유형 Boolean 추가.
  • KIP-884: 색다른 KafkaClientSupplier를 코드 변경 없이 사용할 수 있는 새 구성 default.client.supplier 추가.

3.4.0의 Streams API 변경

  • KIP-770: 구성 cache.max.bytes.buffering 폐기, 새 구성 statestore.cache.max.bytes 사용. 새 메트릭 input-buffer-bytes-total, cache-size-bytes-total 추가.
  • KIP-873: 결과 레코드를 다운스트림 싱크 토픽의 여러 파티션으로 멀티캐스트하고, 보내지 않고 결과 레코드를 버리는 기능 추가. Integer StreamPartitioner.partition() 폐기.
  • KIP-862: 스트림-스트림 셀프 조인용 DSL 최적화(single.store.self.join).
  • KIP-865: 애플리케이션 리셋 도구의 --bootstrap-servers 파라미터 폐기, 새 --bootstrap-server 도입.

3.3.0의 Streams API 변경

  • KIP-812: KafkaStreams.close(CloseOptions) 오버로드 추가로 인스턴스가 즉시 그룹을 떠나도록 강제.
  • KIP-820: PAPI 타입 안전성 개선을 DSL에 적용. KStream.transform, flatTransform, transformValues, flatTransformValuesvoid KStream.process의 오버로드 폐기, 새 KStream.process(ProcessorSupplier, ...)·KStream.processValues(FixedKeyProcessorSupplier, ...) 사용. 주의: processValues()가 회귀 버그(KAFKA-19668)를 도입. "merge repartition topics" 최적화가 활성화되어 있으면 3.3.0에서 transformValues()에서 processValues()로 마이그레이션하는 것은 안전하지 않아요.
  • KIP-825: windowed 집계 결과를 윈도우가 닫힌 후에만 방출하는 "emit strategies" 도입.
  • KIP-834: KafkaStreams.pause(), resume(), isPaused() 추가.
  • KIP-846: 새 토픽 레벨 스코프 메트릭 bytes-consumed-total, records-consumed-total, bytes-produced-total, records-produced-total 추가.

3.2.0의 Streams API 변경

  • KIP-471: RocksDB 메트릭을 다른 카프카 메트릭처럼 접근 가능하게 만드는 것 완성.
  • KIP-591: 모든 DSL 연산자의 기본 저장소를 전역으로 설정하는 새 구성 default.dsl.store 추가.
  • KIP-708: 다중 AZ 배포를 위해 랙 인지 스탠바이 태스크 할당 전략 구성 (rack.aware.assignment.tagsclient.tag.<myTag>).
  • KIP-791: StateStoreContext.recordMetadata() 추가.
  • KIP-796: 완전히 새 IQv2 API 도입 (StateQueryRequest, StateQueryResult, Query, QueryResult). 여러 내장 쿼리 유형 추가.
  • KIP-811: 새 구성 repartition.purge.interval.ms 추가.

3.1.0의 Streams API 변경

  • KIP-633: left/outer 스트림-스트림 조인의 의미론 개선. 조인 결과는 조인 윈도우가 닫힌 후에만 방출. 새 API JoinWindows.ofTimeDifferenceAndGrace().ofTimeDifferenceWithNoGrace(). 윈도우 집계에서도 유예 기간 설정이 필수가 됨.
  • KIP-761: 기본 컨슈머·프로듀서 클라이언트의 차단 시간을 추적하는 새 메트릭.
  • KIP-763·KIP-766: 인터랙티브 쿼리의 범위 쿼리가 null을 개방 경계로 허용.
  • KIP-775: 외래 키 테이블-테이블 조인이 커스텀 파티셔너 지원.

3.0.0의 Streams API 변경

  • KIP-695: 태스크 유휴 상태(max.task.idle.ms) 의미론 개선으로 더 강한 순서 있는 조인·머지 처리 의미론.
  • KIP-216: 인터랙티브 쿼리가 다른 오류에 대해 새 예외(UnknownStateStoreException, StreamsNotStartedException, InvalidStateStorePartitionException)를 던질 수 있음.
  • KIP-732: "exactly_once"(EOS v1) 폐기, 향상된 "exactly_once_v2" 사용. "exactly_once_beta"라는 구성 값도 "exactly_once_v2"로 이름 변경·폐기.
  • 윈도우 연산(윈도우/세션 집계, 스트림-스트림 조인)의 기본 24시간 유예 기간 제거. 모든 정적 생성자를 #ofSizeAndGrace#ofSizeWithNoGrace 유형으로 대체 (Windows 클래스).
  • left/outer 스트림-스트림 조인 결과를 조기에 방출하던 동작 변경(KAFKA-10847) — 가짜 결과 방지.
  • TaskId의 공개 topicGroupId·partition 필드 폐기. TaskMetadata, ThreadMetadata, StreamsMetadata, KafkaStreams#allMetadata 등 폐기 및 새 API 마이그레이션 (KIP-740, KIP-744).
  • replication.factor 기본값을 -1로 변경(브로커 기본 복제 팩터 사용, 브로커 2.4+ 필요).
  • 새 Serde 유형 ListSerde 도입.

2.8.0의 Streams API 변경

  • KIP-689: StreamJoinedwithLoggingEnabled(), withLoggingDisabled() 옵션 추가.
  • KIP-663: KafkaStreams#addStreamThread(), removeStreamThread() 추가.
  • KIP-671: setUncaughtExceptionHandler의 새 시그니처.
  • KIP-696: KafkaStreams 클라이언트 상태 머신 업데이트. ERROR 상태는 이제 터미널이고 PENDING_ERROR가 전환 상태.
  • KIP-572: task.timeout.ms 구성 (기본 5분).

2.7.0의 Streams API 변경

  • KIP-648: KeyQueryMetadata의 getter 이름 변경(getActiveHost()activeHost() 등).
  • KIP-626: StreamsConfig 변수 이름 변경.
  • KIP-450: SlidingWindows를 윈도우 집계 옵션으로 추가.

2.6.0의 Streams API 변경

  • KIP-447: 새 처리 모드 EOS version 2 추가 (processing.guarantee"exactly_once_beta").
  • KIP-441: 고가용성 상태 저장 애플리케이션을 위한 태스크 할당 알고리즘 수정.
  • KIP-221: KStream.through() 폐기, 새 KStream.repartition() 사용.
  • KIP-401: StateStore 사용 편의성 개선.
  • KIP-571: StreamsResetter--force 옵션 추가.
  • KIP-446: Suppressed.withLoggingDisabled()/withLoggingEnabled() 추가.

2.5.0의 Streams API 변경

  • KIP-150: 새 cogroup() 연산자 추가.
  • KIP-523: 새 KStream.toTable() API 추가.
  • KIP-527: 새 Serde 유형 Void 추가.
  • KIP-530: UsePreviousTimeOnInvalidTimestamp 폐기, UsePartitionTimeOnInvalidTimeStamp로 교체.
  • KIP-562·KIP-535: KafkaStreams.store(StoreQueryParameters) 추가.

2.4.0의 Streams API 변경

  • KIP-213: KTable-KTable 외래 키 조인 추가 (INNER와 LEFT 모두).
  • KIP-307: DSL 토폴로지에서 모든 연산자 이름 지정 가능.
  • KIP-479: StreamJoined 클래스 추가로 조인 프로세서·리파티션 토픽·상태 저장소 이름 지정 가능.
  • KIP-429: 증분 협력적 리밸런싱 도입.
  • KIP-444·KIP-471: 새 메트릭 추가.
  • KIP-470: test-utils 개선. TestInputTopic·TestOutputTopic, TestRecord 도입.
  • KIP-528: PartitionGrouper 인터페이스 폐기.

Streams API 브로커 호환성

다음 표는 다양한 카프카 브로커 버전과 호환되는 Kafka Streams API 버전을 보여줘요. 2.4.x보다 오래된 줄의 경우 3.9 업그레이드 문서를 참고해요.

카프카 브로커 (열) / Kafka Streams API (행) 2.4.x - 4.0.x 4.1.x - 4.3.x
2.4.x - 2.5.x 호환 호환
2.6.x - 4.3.x 호환; exactly-once v2 활성화는 브로커 2.5.x 이상 필요 호환

더 알아보기