Streams DSL
Streams DSL
Kafka Streams DSL(Domain Specific Language)은 Streams Processor API 위에 만들어진 고수준 API예요. 대부분의 사용자, 특히 초보자에게 권장돼요. 대부분의 데이터 처리 연산을 몇 줄의 DSL 코드로 표현할 수 있기 때문이에요. 스트림을 만들고, 변환하고, 집계·조인하고, 결과를 다시 카프카로 쓰는 흐름을 이 페이지에서 자세히 풀어드릴게요.
출처: 문서
본문
Kafka Streams DSL은 Streams Processor API 위에 구축돼요. 대부분의 사용자, 특히 초보자에게 권장돼요. 대부분의 데이터 처리 연산은 몇 줄의 DSL 코드로 표현할 수 있어요.
- 집계 (Aggregating)
- 조인 (Joining)
- 조인 공동 파티셔닝 요구사항
- KStream-KStream 조인
- KTable-KTable Equi-Join
- KTable-KTable Foreign-Key 조인
- KStream-KTable 조인
- KStream-GlobalKTable 조인
- 윈도잉 (Windowing)
- Hopping 시간 윈도우
- Tumbling 시간 윈도우
- Sliding 시간 윈도우
- 세션 윈도우
- 윈도우 최종 결과
- 프로세서 적용 (Processor API 통합)
- Transformers 제거 및 프로세서로의 마이그레이션
- Streams DSL 애플리케이션에서 연산자 이름 짓기
- KTable 업데이트율 제어
- 테이블 프로세서에 타임스탬프 기반 의미론 사용
- 스트림을 카프카로 다시 쓰기
- Streams 애플리케이션 테스트
- Kafka Streams DSL for Scala
- 샘플 사용법
- Implicit Serdes
- 사용자 정의 Serdes
개요
Processor API와 비교할 때 오직 DSL만 지원하는 것:
- KStream, KTable, GlobalKTable 형태의 스트림·테이블 내장 추상화. 스트림과 테이블에 대한 일급 지원이 중요한 이유는, 실제로 대부분의 사용 사례가 스트림 또는 데이터베이스/테이블 중 하나가 아니라 둘의 조합을 요구하기 때문이에요. 예를 들어 실시간으로 업데이트되는 고객 360도 뷰를 만들려면 애플리케이션은 많은 고객 관련 이벤트 입력 스트림을, 고객의 지속적으로 업데이트되는 360도 뷰를 포함하는 출력 테이블로 변환하게 돼요.
- 선언적·함수형 프로그래밍 스타일 —
map,filter같은 무상태 변환과, 집계(count,reduce), 조인(leftJoin), 윈도잉(세션 윈도우) 같은 상태 저장 변환을 모두 지원해요.
DSL로 애플리케이션에서 프로세서 토폴로지(즉 논리적 처리 계획)를 정의할 수 있어요. 이를 수행하는 단계:
- 카프카 토픽에서 읽는 하나 이상의 입력 스트림을 지정.
- 이 스트림들에 변환을 구성.
- 결과 출력 스트림을 카프카 토픽으로 다시 쓰거나, 인터랙티브 쿼리를 통해(예: REST API) 애플리케이션의 처리 결과를 다른 애플리케이션에 직접 노출.
애플리케이션 실행 후 정의된 프로세서 토폴로지는 지속적으로 실행돼요(즉 처리 계획이 실행됨). 사용 가능한 API 기능의 전체 목록은 Streams API docs도 참고해요.
KStream
Kafka Streams DSL에만 KStream의 개념이 있어요. KStream은 **레코드 스트림(record stream)**의 추상화이며, 각 데이터 레코드는 무한 데이터 집합에서 자급자족하는 하나의 데이터를 나타내요. 테이블 비유를 쓰면, 레코드 스트림의 데이터 레코드는 항상 "INSERT"로 해석돼요 — 추가 전용 원장(append-only ledger)에 항목을 추가한다고 생각하면 돼요 — 같은 키의 기존 행을 대체하는 레코드는 없기 때문이에요. 예로는 신용카드 거래, 페이지뷰 이벤트, 서버 로그 항목이 있어요.
예를 들어 두 데이터 레코드가 스트림으로 보내진다고 상상해보세요:
("alice", 1) -> ("alice", 3)
스트림 처리 애플리케이션이 사용자별 값을 합한다면 alice에 대해 4를 반환할 거예요. 왜냐하면 두 번째 데이터 레코드가 이전 레코드의 업데이트로 간주되지 않기 때문이에요. 이 KStream 동작을 아래 KTable과 비교해보세요. KTable은 alice에 대해 3을 반환할 거예요.
KTable
Kafka Streams DSL에만 KTable의 개념이 있어요. KTable은 **체인지로그 스트림(changelog stream)**의 추상화이며, 각 데이터 레코드는 업데이트를 나타내요. 더 정확히는 데이터 레코드의 값이 같은 레코드 키의 마지막 값의 "UPDATE"로 해석돼요 (해당 키가 아직 없으면 업데이트는 INSERT로 간주). 테이블 비유를 쓰면, 체인지로그 스트림의 데이터 레코드는 UPSERT(일명 INSERT/UPDATE)로 해석돼요. 같은 키를 가진 기존 행이 덮어쓰이기 때문이에요. 또한 null 값은 특별히 해석돼요: null 값을 가진 레코드는 해당 레코드 키에 대한 "DELETE" 또는 톰스톤(tombstone)을 나타내요.
예를 들어 두 데이터 레코드가 스트림으로 보내진다고 상상해보세요:
("alice", 1) -> ("alice", 3)
스트림 처리 애플리케이션이 사용자별 값을 합한다면 alice에 대해 3을 반환할 거예요. 두 번째 데이터 레코드가 이전 레코드의 업데이트로 간주되기 때문이에요.
카프카 로그 컴팩션의 효과: KStream과 KTable을 생각하는 또 다른 방법: KTable을 카프카 토픽에 저장한다면 저장 공간을 아끼기 위해 카프카의 로그 컴팩션 기능을 활성화하고 싶을 거예요. 그러나 KStream의 경우 로그 컴팩션을 활성화하는 것은 안전하지 않아요. 로그 컴팩션이 같은 키의 이전 데이터 레코드를 제거하기 시작하는 순간 데이터의 의미론이 깨지기 때문이에요. 예시를 다시 들면, 로그 컴팩션이 ("alice", 1) 데이터 레코드를 제거했기 때문에 갑자기 alice에 대해 4 대신 3을 얻을 거예요. 따라서 로그 컴팩션은 KTable(체인지로그 스트림)에는 완전히 안전하지만 KStream(레코드 스트림)에는 위험해요.
체인지로그 스트림의 예는 스트림·테이블 섹션에서 이미 보았어요. 또 다른 예는 관계형 데이터베이스의 체인지로그에 있는 CDC(change data capture) 레코드로, 데이터베이스 테이블의 어느 행이 삽입, 업데이트, 삭제되었는지를 나타내요.
KTable은 키로 데이터 레코드의 현재 값을 조회하는 기능도 제공해요. 이 테이블-조회 기능은 조인 연산(개발자 가이드의 Joining 참고)과 인터랙티브 쿼리를 통해 사용할 수 있어요.
GlobalKTable
Kafka Streams DSL에만 GlobalKTable의 개념이 있어요. KTable처럼 GlobalKTable은 체인지로그 스트림의 추상화이며, 각 데이터 레코드는 업데이트를 나타내요.
GlobalKTable이 KTable과 다른 점은 채워지는 데이터, 즉 기본 카프카 토픽의 어떤 데이터가 각 테이블로 읽히는지예요. 약간 단순화해서, 5개의 파티션이 있는 입력 토픽이 있다고 상상해보세요. 애플리케이션에서 이 토픽을 테이블로 읽고 싶어요. 또한 최대 병렬성을 위해 5개의 애플리케이션 인스턴스에서 애플리케이션을 실행하고 싶어요.
- 입력 토픽을 KTable로 읽으면 각 애플리케이션 인스턴스의 "로컬" KTable 인스턴스는 토픽의 5개 파티션 중 1개 파티션의 데이터로만 채워져요.
- 입력 토픽을 GlobalKTable로 읽으면 각 애플리케이션 인스턴스의 로컬 GlobalKTable 인스턴스는 토픽의 모든 파티션의 데이터로 채워져요.
GlobalKTable은 키로 데이터 레코드의 현재 값을 조회하는 기능을 제공해요. 이 테이블-조회 기능은 join operations을 통해 사용할 수 있어요. GlobalKTable은 KTable과 달리 시간의 개념이 없다는 점을 유의해요.
글로벌 테이블의 이점:
- 더 편리하고/효율적인 조인: 특히 글로벌 테이블은 스타 조인(star join)을 수행할 수 있고, "외래 키" 조회(레코드 키뿐 아니라 레코드 값의 데이터로도 테이블 데이터를 조회)를 지원하며, 여러 조인을 연결할 때 더 효율적이에요. 또한 글로벌 테이블과 조인할 때 입력 데이터는 공동 파티셔닝(co-partitioning)이 필요 없어요.
- 애플리케이션의 모든 실행 인스턴스에 정보를 "브로드캐스트"하는 데 사용할 수 있어요.
글로벌 테이블의 단점:
- (파티셔닝된) KTable보다 로컬 저장소 소비가 증가해요. 전체 토픽을 추적하기 때문이에요.
- (파티셔닝된) KTable보다 네트워크와 카프카 브로커 부하가 증가해요. 전체 토픽을 읽기 때문이에요.
카프카에서 소스 스트림 만들기
카프카 토픽에서 애플리케이션으로 데이터를 쉽게 읽을 수 있어요. 다음 연산이 지원돼요.
| 카프카에서 읽기 | 설명 |
|---|---|
Stream input topics -> KStream |
지정된 카프카 입력 토픽에서 KStream을 만들고 데이터를 레코드 스트림으로 해석. KStream은 파티셔닝된 레코드 스트림을 나타냄. KStream의 경우 모든 애플리케이션 인스턴스의 로컬 KStream 인스턴스는 입력 토픽 파티션의 일부에서만 데이터로 채워져요. 모든 애플리케이션 인스턴스에 걸쳐 모든 입력 토픽 파티션이 읽히고 처리돼요. |
Table input topic -> KTable |
지정된 카프카 입력 토픽을 KTable로 읽음. 토픽은 체인지로그 스트림으로 해석되며, 같은 키를 가진 레코드는 UPSERT(일명 INSERT/UPDATE, 값이 null이 아닐 때) 또는 해당 키에 대한 DELETE(값이 null일 때)로 해석돼요. 테이블(더 정확히는 테이블을 뒷받침하는 내부 상태 저장소)에 이름을 제공해야 해요. 이것은 테이블에 대한 인터랙티브 쿼리를 지원하는 데 필요해요. 이름을 제공하지 않으면 테이블은 쿼리 가능하지 않고 상태 저장소에는 내부 이름이 제공돼요. |
Global Table input topic -> GlobalKTable |
지정된 카프카 입력 토픽을 GlobalKTable로 읽음. 토픽은 체인지로그 스트림으로 해석되며, 같은 키를 가진 레코드는 UPSERT(값이 null이 아닐 때) 또는 DELETE(값이 null일 때)로 해석돼요. GlobalKTable의 경우 모든 애플리케이션 인스턴스의 로컬 GlobalKTable 인스턴스는 입력 토픽의 모든 파티션 데이터로 채워져요. 테이블(내부 상태 저장소)에 이름을 제공해야 해요. |
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.KStream;
StreamsBuilder builder = new StreamsBuilder();
KStream<String, Long> wordCounts = builder.stream(
"word-counts-input-topic", /* input topic */
Consumed.with(
Serdes.String(), /* key serde */
Serdes.Long() /* value serde */
);
Serdes를 명시적으로 지정하지 않으면 구성의 기본 Serdes가 사용돼요. 카프카 입력 토픽의 레코드 키나 값 유형이 구성된 기본 Serdes와 일치하지 않으면 Serdes를 명시적으로 지정해야 해요. 기본 Serdes 구성, 사용 가능한 Serdes, 커스텀 Serdes 구현에 대한 정보는 데이터 타입과 직렬화를 참고해요.
stream의 여러 변형이 있어요. 예를 들어 읽을 입력 토픽의 정규식 패턴을 지정할 수 있어요 (이렇게 구독하면 일치하는 모든 토픽이 같은 입력 토픽 그룹의 일부가 되고, 다른 토픽에 대해 작업이 병렬화되지 않는다는 점을 유의하세요).
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.GlobalKTable;
StreamsBuilder builder = new StreamsBuilder();
GlobalKTable<String, Long> wordCounts = builder.globalTable(
"word-counts-input-topic",
Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as(
"word-counts-global-store" /* table/store name */)
.withKeySerde(Serdes.String()) /* key serde */
.withValueSerde(Serdes.Long()) /* value serde */
);
카프카 입력 토픽의 레코드 키나 값 유형이 구성된 기본 Serdes와 일치하지 않으면 Serdes를 명시적으로 지정해야 해요.
스트림 변환
KStream과 KTable 인터페이스는 다양한 변환 연산을 지원해요. 각 연산은 하나 이상의 연결된 프로세서로 변환되어 기본 프로세서 토폴로지가 돼요. KStream과 KTable은 강하게 타입화 되어 있으므로, 이 변환 연산들은 모두 사용자가 입력·출력 데이터 타입을 지정할 수 있는 제네릭 함수로 정의돼요.
일부 KStream 변환은 하나 이상의 KStream 객체를 생성할 수 있어요. 예: KStream의 filter와 map은 다른 KStream을 생성하고, KStream의 split은 여러 KStream을 생성할 수 있어요. 일부 다른 변환은 KTable 객체를 생성할 수 있어요. 예: KStream의 집계도 KTable을 생성해요. 이것은 Kafka Streams가 이미 다운스트림 변환 연산자에 생성된 후 순서가 뒤집힌 레코드가 도착하면 계산된 값을 지속적으로 업데이트할 수 있게 해줘요.
모든 KTable 변환 연산은 다른 KTable만 생성할 수 있어요. 그러나 Kafka Streams DSL은 KTable 표현을 KStream으로 변환하는 특수 함수를 제공해요. 이 변환 메서드는 모두 연결되어 복잡한 프로세서 토폴로지를 구성할 수 있어요.
변환 연산은 두 가지 하위 섹션으로 설명돼요:
- 무상태 변환 (Stateless transformations)
- 상태 저장 변환 (Stateful transformations)
무상태 변환
무상태 변환은 처리를 위해 상태가 필요 없고 스트림 프로세서와 연결된 상태 저장소도 필요 없어요. 카프카 0.11.0 이상에서는 무상태 KTable 변환의 결과를 구체화(materialize)할 수 있어요. 이것은 인터랙티브 쿼리를 통해 결과를 쿼리할 수 있게 해줘요. KTable을 구체화하려면 아래 각 무상태 연산에 선택적 queryableStoreName 인자를 추가할 수 있어요.
| 변환 | 설명 |
|---|---|
Branch KStream -> BranchedKStream |
공급된 프레디킷에 따라 KStream을 하나 이상의 KStream 인스턴스로 분기(또는 분할). 프레디킷은 순서대로 평가됨. 레코드는 첫 번째 일치에서 하나의 출력 스트림에만 배치됨: n번째 프레디킷이 true이면 레코드는 n번째 스트림에 배치됨. 어떤 프레디킷에도 일치하지 않으면 기본 분기로 라우팅되거나, 기본 분기가 없으면 버려짐. 분기는 예를 들어 레코드를 다른 다운스트림 토픽으로 라우팅하는 데 유용. |
| Broadcast/Multicast | KStream을 여러 다운스트림 연산자로 브로드캐스트. 같은 KStream 인스턴스에 여러 연산자를 적용해 레코드를 둘 이상의 연산자로 보냄. 멀티캐스트는 브로드캐스트 + 필터로 구현됨. 각 레코드를 최대 하나의 다운스트림 분기로 보내는 분기와 달리, 멀티캐스트는 레코드를 원하는 수만큼의 다운스트림 KStream 인스턴스로 보낼 수 있음. |
Filter KStream -> KStream, KTable -> KTable |
각 요소에 대해 boolean 함수를 평가하고 함수가 true를 반환하는 요소를 유지. |
| Inverse Filter | 각 요소에 대해 boolean 함수를 평가하고 함수가 true를 반환하는 요소를 버림 (filterNot). |
FlatMap KStream -> KStream |
하나의 레코드를 받아 0, 1 또는 그 이상의 레코드를 생성. 레코드 키와 값을 포함해 타입을 수정할 수 있음. 스트림을 데이터 재파티셔닝으로 표시: flatMap 후 그룹화·조인을 적용하면 레코드가 재파티셔닝됨. 가능하면 재파티셔닝을 일으키지 않는 flatMapValues 사용. |
FlatMapValues KStream -> KStream |
하나의 레코드를 받아 0, 1 또는 그 이상의 레코드를 생성, 원래 레코드의 키를 유지. 값과 값 타입을 수정할 수 있음. flatMap보다 선호됨 — 재파티셔닝을 일으키지 않기 때문. 그러나 flatMap처럼 키나 키 타입을 수정할 수 없음. |
Foreach KStream -> void |
종단(terminal) 연산. 각 레코드에 무상태 액션을 수행. 부작용을 일으키는 데 사용(peek과 유사)하고 이후 입력 데이터의 추가 처리를 중지(peek과 달리, peek은 종단 연산이 아님). 처리 보장 참고: 액션의 부작용(외부 시스템 쓰기 같은)은 카프카가 추적할 수 없어 보통 카프카의 처리 보장 혜택을 받지 못함. |
GroupByKey KStream -> KGroupedStream |
기존 키로 레코드를 그룹화. 그룹화는 스트림·테이블 집계의 전제 조건이며 후속 연산을 위해 데이터가 제대로 파티셔닝("키잉")되도록 보장. 스트림이 재파티셔닝으로 표시된 경우에만 데이터 재파티셔닝을 일으킴. groupBy보다 선호됨 — 스트림이 이미 재파티셔닝으로 표시된 경우에만 데이터를 재파티셔닝하기 때문. 그러나 groupBy처럼 키나 키 타입을 수정할 수 없음. |
GroupBy KStream -> KGroupedStream, KTable -> KGroupedTable |
새 키(다른 키 타입일 수 있음)로 레코드를 그룹화. 테이블을 그룹화할 때 새 값과 값 타입을 지정할 수도 있음. selectKey(...).groupByKey()의 축약형. 항상 데이터 재파티셔닝을 일으킴. 가능하면 groupByKey 사용. |
Cogroup KGroupedStream -> CogroupedKStream |
여러 입력 스트림을 단일 연산으로 집계할 수 있게 함. (이미 그룹화된) 다른 입력 스트림은 같은 키 타입을 가져야 하고 다른 값 타입을 가질 수 있음. KGroupedStream#cogroup()은 단일 입력 스트림으로 새 cogrouped 스트림을 만들고, CogroupedKStream#cogroup()은 기존 cogrouped 스트림에 그룹화된 스트림을 추가. CogroupedKStream은 집계되기 전에 윈도윙될 수 있음. |
Map KStream -> KStream |
하나의 레코드를 받아 하나의 레코드를 생성. 레코드 키와 값을 포함해 타입을 수정할 수 있음. 스트림을 재파티셔닝으로 표시. 가능하면 mapValues 사용. |
Map (values only) KStream -> KStream, KTable -> KTable |
하나의 레코드를 받아 하나의 레코드를 생성, 원래 레코드의 키를 유지. 값과 값 타입을 수정할 수 있음. map보다 선호됨 — 재파티셔닝을 일으키지 않기 때문. 그러나 map처럼 키나 키 타입을 수정할 수 없음. |
Merge KStream -> KStream |
두 스트림의 레코드를 하나의 더 큰 스트림으로 병합. 병합된 스트림에서 서로 다른 스트림의 레코드 사이에 순서 보장은 없음. 각 입력 스트림 내에서는 상대 순서가 보존됨. |
Peek KStream -> KStream |
각 레코드에 무상태 액션을 수행하고 변경되지 않은 스트림을 반환. 부작용을 일으키는 데 사용하고 입력 데이터 처리를 계속. foreach와 달리 종단 연산이 아님. peek은 입력 스트림을 그대로 반환. 로깅, 메트릭 추적, 디버깅·문제 해결에 유용. |
Print KStream -> void |
종단 연산. 레코드를 System.out으로 출력. print() 호출은 foreach((key, value) -> System.out.println(key + ", " + value))와 같음. 주로 디버깅/테스트용이며 각 레코드 출력 시 플러시를 시도함. 성능 요구사항이 중요하면 프로덕션에서 사용하지 말 것. |
SelectKey KStream -> KStream |
각 레코드에 새 키(새 키 타입일 수 있음)를 할당. selectKey(mapper)는 map((key, value) -> mapper(key, value), value)와 같음. 스트림을 재파티셔닝으로 표시. |
Table to Stream KTable -> KStream |
이 테이블의 체인지로그 스트림을 얻음. toStream의 변형으로 결과 스트림에 새 키를 선택할 수 있는 것도 있음. |
Stream to Table KStream -> KTable |
이벤트 스트림을 테이블(체인지로그 스트림)로 변환 (toTable). |
Repartition KStream -> KStream |
원하는 파티션 수로 스트림의 재파티셔닝을 수동으로 트리거. repartition()은 항상 스트림의 재파티셔닝을 트리거하므로, 키 변경 연산을 미리 수행해도 자동 재파티셔닝을 트리거하지 않는 임베드된 Processor API 메서드(process() 등)와 함께 사용할 수 있음. |
// Branch 예시
KStream<String, Long> stream = ...;
Map<String, KStream<String, Long>> branches =
stream.split(Named.as("Branch-"))
.branch((key, value) -> key.startsWith("A"), /* first predicate */
Branched.as("A"))
.branch((key, value) -> key.startsWith("B"), /* second predicate */
Branched.as("B"))
.defaultBranch(Branched.as("C")) /* default branch */
);
// KStream branches.get("Branch-A") contains all records whose keys start with "A"
// KStream branches.get("Branch-B") contains all records whose keys start with "B"
// KStream branches.get("Branch-C") contains all other records
// Filter 예시
KStream<String, Long> stream = ...;
// A filter that selects (keeps) only positive numbers
KStream<String, Long> onlyPositives = stream.filter((key, value) -> value > 0);
// FlatMap 예시
KStream<Long, String> stream = ...;
KStream<String, Integer> transformed = stream.flatMap(
// Here, we generate two output records for each input record.
// We also change the key and value types.
// Example: (345L, "Hello") -> ("HELLO", 1000), ("hello", 9000)
(key, value) -> {
List<KeyValue<String, Integer>> result = new LinkedList<>();
result.add(KeyValue.pair(value.toUpperCase(), 1000));
result.add(KeyValue.pair(value.toLowerCase(), 9000));
return result;
}
);
// FlatMapValues 예시 — 문장을 단어로 쪼개기
KStream<byte[], String> sentences = ...;
KStream<byte[], String> words = sentences.flatMapValues(value -> Arrays.asList(value.split("\\s+")));
// Foreach 예시
KStream<String, Long> stream = ...;
// Print the contents of the KStream to the local console.
stream.foreach((key, value) -> System.out.println(key + " => " + value));
// GroupByKey 예시
KStream<byte[], String> stream = ...;
// Group by the existing key, using the application's configured
// default serdes for keys and values.
KGroupedStream<byte[], String> groupedStream = stream.groupByKey();
// When the key and/or value types do not match the configured
// default serdes, we must explicitly specify serdes.
KGroupedStream<byte[], String> groupedStream = stream.groupByKey(
Grouped.with(
Serdes.ByteArray(), /* key */
Serdes.String()) /* value */
);
// GroupBy 예시
KStream<byte[], String> stream = ...;
KTable<byte[], String> table = ...;
// Group the stream by a new key and key type
KGroupedStream<String, String> groupedStream = stream.groupBy(
(key, value) -> value,
Grouped.with(
Serdes.String(), /* key (note: type was modified) */
Serdes.String()) /* value */
);
// Group the table by a new key and key type, and also modify the value and value type.
KGroupedTable<String, Integer> groupedTable = table.groupBy(
(key, value) -> KeyValue.pair(value, value.length()),
Grouped.with(
Serdes.String(), /* key (note: type was modified) */
Serdes.Integer()) /* value (note: type was modified) */
);
// Cogroup 예시
KStream<byte[], String> stream = ...;
KStream<byte[], String> stream2 = ...;
// Group by the existing key, using the application's configured
// default serdes for keys and values.
KGroupedStream<byte[], String> groupedStream = stream.groupByKey();
KGroupedStream<byte[], String> groupedStream2 = stream2.groupByKey();
CogroupedKStream<byte[], String> cogroupedStream = groupedStream.cogroup(aggregator1).cogroup(groupedStream2, aggregator2);
KTable<byte[], String> table = cogroupedStream.aggregate(initializer);
KTable<byte[], String> table2 = cogroupedStream.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMillis(500))).aggregate(initializer);
// Map 예시
KStream<byte[], String> stream = ...;
// Note how we change the key and the key type (similar to `selectKey`)
// as well as the value and the value type.
KStream<String, Integer> transformed = stream.map(
(key, value) -> KeyValue.pair(value.toLowerCase(), value.length()));
// mapValues 예시
KStream<byte[], String> stream = ...;
KStream<byte[], String> uppercased = stream.mapValues(value -> value.toUpperCase());
// Merge 예시
KStream<byte[], String> stream1 = ...;
KStream<byte[], String> stream2 = ...;
KStream<byte[], String> merged = stream1.merge(stream2);
// Print 예시
KStream<byte[], String> stream = ...;
// print to sysout
stream.print();
// print to file with a custom label
stream.print(Printed.toFile("streams.out").withLabel("streams"));
// SelectKey 예시
KStream<byte[], String> stream = ...;
// Derive a new record key from the record's value. Note how the key type changes, too.
KStream<String, String> rekeyed = stream.selectKey((key, value) -> value.split(" ")[0]);
// Table to Stream 예시
KTable<byte[], String> table = ...;
// Also, a variant of `toStream` exists that allows you
// to select a new key for the resulting stream.
KStream<byte[], String> stream = table.toStream();
// Repartition 예시
KStream<byte[], String> stream = ... ;
KStream<byte[], String> repartitionedStream = stream.repartition(Repartitioned.numberOfPartitions(10));
상태 저장 변환
상태 저장 변환은 입력을 처리하고 출력을 생성하는 데 상태에 의존하며, 스트림 프로세서와 연결된 상태 저장소가 필요해요. 예를 들어 집계 연산에서 윈도잉 상태 저장소가 윈도우별 최신 집계 결과를 수집하는 데 사용돼요. 조인 연산에서 윈도잉 상태 저장소가 정의된 윈도우 경계 내에서 지금까지 받은 모든 레코드를 수집하는 데 사용돼요.
다음 저장소 유형이 (파라미터 materialized로 지정된 유형과 무관하게) 사용돼요:
- 비윈도우 집계와 비윈도우 KTable은
TimestampedKeyValueStores또는VersionedKeyValueStores를 사용 —materialized파라미터가 versioned인지에 따라. - 시간-윈도우 집계와 KStream-KStream 조인은
TimestampedWindowStores를 사용. - 세션 윈도우 집계는
SessionStores를 사용 (현재 timestamped 세션 저장소는 없음).
헤더 인지 상태 저장소 (KIP-1285): dsl.store.format=HEADERS로 설정하면 지원되는 DSL 연산자가 헤더 인지 상태 저장소를 사용하게 해요. 이 저장소들은 레코드 헤더를 값·타임스탬프와 함께 유지할 수 있어요. 이 구성은 상태 저장소 형식만 변경해요. DSL 연산자가 출력 레코드용 헤더를 어떻게 만드는지는 정의하지 않아요. 현재 동작은:
- 집계, KTable-KTable 조인, 구체화된
KTable.mapValues,KStream.toTable(),StreamsBuilder.table()은 그 구체화된 저장소에 빈 헤더를 써요. - KStream-KStream 조인 윈도우 저장소는 소스 레코드 헤더를 유지하지만, 조인 결과 레코드는 계산되거나 병합된 헤더를 가지지 않아요; 결과를 트리거한 레코드의 헤더를 가질 수 있어요.
suppress()와 left/outer KStream-KStream 조인은 헤더 인지가 아닌 버퍼 저장소를 사용하므로 그 버퍼를 통과하는 레코드는 헤더를 잃어요.
후속 KIP가 DSL 결과 헤더를 어떻게 계산할지 정의할 거예요.
상태 저장소는 장애 허용적이라는 점을 유의해요. 장애가 발생하면 Kafka Streams는 처리를 재개하기 전에 모든 상태 저장소를 완전히 복원하도록 보장해요 (Fault Tolerance 참고).
DSL에서 사용 가능한 상태 저장 변환:
- 집계 (Aggregating)
- 조인 (Joining)
- 윈도잉 (Windowing) — 집계와 조인의 일부로
- 커스텀 프로세서·트랜스포머 적용 — 상태 저장일 수 있으며 Processor API 통합용
WordCount 예시:
// Assume the record values represent lines of text. For the sake of this example, you can ignore
// whatever may be stored in the record keys.
KStream<String, String> textLines = ...;
KStream<String, Long> wordCounts = textLines
// Split each text line, by whitespace, into words. The text lines are the record
// values, i.e. you can ignore whatever data is in the record keys and thus invoke
// `flatMapValues` instead of the more generic `flatMap`.
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
// Group the stream by word to ensure the key of the record is the word.
.groupBy((key, word) -> word)
// Count the occurrences of each word (record key).
//
// This will change the stream type from `KGroupedStream<String, String>` to
// `KTable<String, Long>` (word -> count).
.count()
// Convert the `KTable<String, Long>` into a `KStream<String, Long>`.
.toStream();
집계 (Aggregating)
레코드가 groupByKey 또는 groupBy로 키별 그룹화된 후 — 따라서 KGroupedStream 또는 KGroupedTable로 표현됨 — reduce 같은 연산으로 집계될 수 있어요. 집계는 키 기반 연산이에요, 즉 항상 같은 키의 레코드(특히 레코드 값)에 대해 동작해요. 윈도우되거나 윈도우되지 않은 데이터에 집계를 수행할 수 있어요.
| 변환 | 설명 |
|---|---|
Aggregate KGroupedStream -> KTable, KGroupedTable -> KTable |
롤링 집계. 그룹화된 키(또는 cogroup)에 의해 (비윈도우) 레코드의 값을 집계. 집계는 reduce의 일반화이며 집계 값이 입력 값과 다른 타입을 가질 수 있게 해줌. 그룹화된 스트림을 집계할 때는 initializer(예: aggValue = 0)와 "adder" 집계자(예: aggValue + curValue)를 제공해야 함. 그룹화된 테이블을 집계할 때는 추가로 "subtractor" 집계자(예: aggValue - oldValue)를 제공해야 함. cogrouped 스트림을 집계할 때는 실제 집계자가 이전 cogroup() 호출에서 각 입력 스트림에 제공되므로 initializer만 제공하면 됨. |
Aggregate (windowed) KGroupedStream -> KTable |
윈도우 집계. 그룹화된 키에 의해 윈도우별로 레코드의 값을 집계. initializer, "adder" 집계자, 윈도우를 제공해야 함. 세션 윈도우를 기반으로 윈도윙할 때는 추가로 "session merger" 집계자(예: mergedAggValue = leftAggValue + rightAggValue)를 제공해야 함. |
Count KGroupedStream -> KTable, KGroupedTable -> KTable |
롤링 집계. 그룹화된 키로 레코드 수를 셈. |
Count (windowed) KGroupedStream -> KTable |
윈도우 집계. 그룹화된 키로 윈도우별 레코드 수를 셈. |
Reduce KGroupedStream -> KTable, KGroupedTable -> KTable |
롤링 집계. 그룹화된 키로 (비윈도우) 레코드의 값을 결합. 현재 레코드 값은 마지막 감소된 값과 결합되고 새 감소된 값을 반환. aggregate와 달리 결과 값 타입을 변경할 수 없음. 그룹화된 스트림을 감소할 때 "adder" reducer(예: aggValue + curValue)를 제공해야 함. 테이블은 추가로 "subtractor" reducer(예: aggValue - oldValue)를 제공해야 함. |
Reduce (windowed) KGroupedStream -> KTable |
윈도우 집계. 그룹화된 키로 윈도우별 레코드 값을 결합. |
// Aggregate (KGroupedStream) 예시
KGroupedStream<byte[], String> groupedStream = ...;
KGroupedTable<byte[], String> groupedTable = ...;
// Aggregating a KGroupedStream (note how the value type changes from String to Long)
KTable<byte[], Long> aggregatedStream = groupedStream.aggregate(
() -> 0L, /* initializer */
(aggKey, newValue, aggValue) -> aggValue + newValue.length(), /* adder */
Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("aggregated-stream-store") /* state store name */
.withValueSerde(Serdes.Long()); /* serde for aggregate value */
// Aggregating a KGroupedTable (note how the value type changes from String to Long)
KTable<byte[], Long> aggregatedTable = groupedTable.aggregate(
() -> 0L, /* initializer */
(aggKey, newValue, aggValue) -> aggValue + newValue.length(), /* adder */
(aggKey, oldValue, aggValue) -> aggValue - oldValue.length(), /* subtractor */
Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("aggregated-table-store") /* state store name */
.withValueSerde(Serdes.Long()) /* serde for aggregate value */
KGroupedStream의 상세 동작:
null키를 가진 입력 레코드는 무시돼요.- 레코드 키가 처음 수신되면 initializer가 호출되고 (adder보다 먼저 호출).
null이 아닌 값의 레코드가 수신될 때마다 adder가 호출돼요.
KGroupedTable의 상세 동작:
null키를 가진 입력 레코드는 무시돼요.- 레코드 키가 처음 수신되면 initializer가 호출되고 (adder·subtractor보다 먼저).
KGroupedStream과 달리, 시간이 지나면서 키에 대한 입력 톰스톤 레코드를 받은 결과로 키에 대해 initializer가 여러 번 호출될 수 있음. - 키에 대해 첫
null아닌 값이 수신되면(예: INSERT), adder만 호출돼요. - 키에 대해 후속
null아닌 값이 수신되면(예: UPDATE), (1) subtractor가 테이블에 저장된 이전 값으로 호출되고 (2) 방금 받은 입력 레코드의 새 값으로 adder가 호출돼요. subtractor는 이전·새 값의 추출된 그룹화 키가 같으면 adder보다 먼저 호출되는 것이 보장돼요. 이 경우 감지는 추출된 키 타입의 equals() 메서드 올바른 구현에 달려 있어요. 그렇지 않으면 subtractor와 adder의 실행 순서는 정의되지 않아요. - 키에 대해 톰스톤 레코드(즉
null값의 레코드; DELETE)가 수신되면 subtractor만 호출돼요. subtractor 자체가null값을 반환하면 해당 키는 결과KTable에서 제거돼요. 그런 일이 발생하면 그 키에 대한 다음 입력 레코드가 initializer를 다시 트리거해요.
// Aggregate (windowed) 예시
import java.time.Duration;
KGroupedStream<String, Long> groupedStream = ...;
// Aggregating with time-based windowing (here: with 5-minute tumbling windows)
KTable<Windowed<String>, Long> timeWindowedAggregatedStream = groupedStream.windowedBy(Duration.ofMinutes(5))
.aggregate(
() -> 0L, /* initializer */
(aggKey, newValue, aggValue) -> aggValue + newValue, /* adder */
Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("time-windowed-aggregated-stream-store") /* state store name */
.withValueSerde(Serdes.Long())); /* serde for aggregate value */
// Aggregating with session-based windowing (here: with an inactivity gap of 5 minutes)
KTable<Windowed<String>, Long> sessionizedAggregatedStream = groupedStream.windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(5)))
.aggregate(
() -> 0L, /* initializer */
(aggKey, newValue, aggValue) -> aggValue + newValue, /* adder */
(aggKey, leftAggValue, rightAggValue) -> leftAggValue + rightAggValue, /* session merger */
Materialized.<String, Long, SessionStore<Bytes, byte[]>>as("sessionized-aggregated-stream-store") /* state store name */
.withValueSerde(Serdes.Long())); /* serde for aggregate value */
집계 의미론 예시 (KGroupedStream -> KTable):
// Key: word, value: count
KStream<String, Integer> wordCounts = ...;
KGroupedStream<String, Integer> groupedStream = wordCounts
.groupByKey(Grouped.with(Serdes.String(), Serdes.Integer()));
KTable<String, Integer> aggregated = groupedStream.aggregate(
() -> 0, /* initializer */
(aggKey, newValue, aggValue) -> aggValue + newValue, /* adder */
Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("aggregated-stream-store" /* state store name */)
.withKeySerde(Serdes.String()) /* key serde */
.withValueSerde(Serdes.Integer()); /* serde for aggregate value */
| Timestamp | 입력 레코드 | 그룹화 | Initializer | Adder | 상태 |
|---|---|---|---|---|---|
| 1 | (hello, 1) | (hello, 1) | 0 (hello용) | (hello, 0 + 1) | (hello, 1) |
| 2 | (kafka, 1) | (kafka, 1) | 0 (kafka용) | (kafka, 0 + 1) | (hello, 1)(kafka, 1) |
| 3 | (streams, 1) | (streams, 1) | 0 (streams용) | (streams, 0 + 1) | (hello, 1)(kafka, 1)(streams, 1) |
| 4 | (kafka, 1) | (kafka, 1) | (kafka, 1 + 1) | (hello, 1)(kafka, 2)(streams, 1) | |
| 5 | (kafka, 1) | (kafka, 1) | (kafka, 2 + 1) | (hello, 1)(kafka, 3)(streams, 1) | |
| 6 | (streams, 1) | (streams, 1) | (streams, 1 + 1) | (hello, 1)(kafka, 3)(streams, 2) |
참고 — 레코드 캐시의 영향:
KTable aggregated열은 테이블의 상태 변화를 시간에 따라 매우 세밀하게 보여줘요. 실제로는 레코드 캐시가 비활성화된 경우에만 그렇게 세밀한 상태 변화를 관찰할 거예요 (기본: 활성화). 레코드 캐시가 활성화되면 timestamps 4·5의 행 출력이 컴팩션되어 KTable에서 키kafka에 대한 단일 상태 업데이트만 있을 수 있어요. 보통 테스트·디버깅 목적으로만 레코드 캐시를 비활성화해야 해요 — 정상적인 상황에서는 활성화 상태로 두는 것이 좋아요.
조인 (Joining)
스트림과 테이블도 조인될 수 있어요. 실제로 많은 스트림 처리 애플리케이션이 스트리밍 조인으로 코딩돼요. 예를 들어 온라인 상점을 뒷받침하는 애플리케이션은 새 데이터 레코드(고객 거래)에 맥락 정보를 강화(enrich)하기 위해 여러 업데이트되는 데이터베이스 테이블(판매 가격, 재고, 고객 정보)에 접근해야 할 수 있어요. 즉 매우 큰 규모로 낮은 처리 지연으로 테이블 조회를 수행해야 하는 시나리오예요. 여기서 유명한 패턴은 카프카의 Connect API와 결합된 소위 CDC(change data capture)를 통해 데이터베이스의 정보를 카프카에서 사용할 수 있게 하고, 각 레코드에 대해 네트워크로 원격 데이터베이스에 쿼리하도록 요구하는 대신 Streams API를 활용해 그러한 테이블·스트림의 매우 빠르고 효율적인 로컬 조인을 수행하는 애플리케이션을 구현하는 것이에요. 이 예시에서 Kafka Streams의 KTable 개념은 각 테이블의 최신 상태(스냅샷)를 로컬 상태 저장소에서 추적할 수 있게 해주며, 따라서 그러한 스트리밍 조인을 할 때 처리 지연을 크게 줄이고 원격 데이터베이스의 부하도 줄여요.
지원되는 조인 연산 (피연산자에 따라 윈도우 조인 또는 비윈도우 조인):
| 조인 피연산자 | 유형 | (INNER) JOIN | LEFT JOIN | OUTER JOIN |
|---|---|---|---|---|
| KStream-to-KStream | 윈도우 | 지원 | 지원 | 지원 |
| KTable-to-KTable | 비윈도우 | 지원 | 지원 | 지원 |
| KTable-to-KTable Foreign-Key Join | 비윈도우 | 지원 | 지원 | 지원 안 함 |
| KStream-to-KTable | 비윈도우 | 지원 | 지원 | 지원 안 함 |
| KStream-to-GlobalKTable | 비윈도우 | 지원 | 지원 | 지원 안 함 |
| KTable-to-GlobalKTable | N/A | 지원 안 함 | 지원 안 함 | 지원 안 함 |
조인 공동 파티셔닝(co-partitioning) 요구사항
equi-조인의 경우 입력 데이터는 조인할 때 공동 파티셔닝되어야 해요. 이것은 조인 양쪽의 같은 키를 가진 입력 레코드가 처리 중에 같은 스트림 태스크로 전달되도록 보장해요. 조인할 때 데이터 공동 파티셔닝을 보장하는 것은 사용자의 책임이에요. KTable-KTable Foreign-Key 조인과 Global KTable 조인을 수행할 때는 공동 파티셔닝이 필요하지 않아요.
데이터 공동 파티셔닝 요구사항:
- 조인의 입력 토픽(왼쪽과 오른쪽)은 같은 수의 파티션을 가져야 해요.
- 입력 토픽에 쓰는 모든 애플리케이션은 같은 파티셔닝 전략을 가져야 해서 같은 키의 레코드가 같은 파티션 번호로 전달돼야 해요. 다시 말해, 입력 데이터의 키스페이스가 같은 방식으로 파티션 전체에 분산되어야 해요. 예를 들어 카프카 Java Producer API를 사용하는 애플리케이션은 같은 파티셔너(생산자 설정
"partitioner.class"일명ProducerConfig.PARTITIONER_CLASS_CONFIG)를 사용해야 하고, 카프카 Streams API를 사용하는 애플리케이션은KStream#to()같은 연산에 대해 같은StreamPartitioner를 사용해야 해요. 모든 애플리케이션에서 기본 파티셔너 관련 설정을 사용한다면 파티셔닝 전략에 대해 걱정할 필요가 없다는 것이 좋은 소식이에요.
데이터 공동 파티셔닝이 왜 필요한가요? KStream-KStream, KTable-KTable, KStream-KTable 조인이 레코드의 키(leftRecord.key == rightRecord.key)를 기반으로 수행되기 때문에, 조인의 입력 스트림/테이블이 키로 공동 파티셔닝되어야 해요.
공동 파티셔닝이 필요하지 않은 두 가지 예외가 있어요. KStream-GlobalKTable 조인의 경우 GlobalKTable의 기본 체인지로그 스트림의 모든 파티션이 각 KafkaStreams 인스턴스에 제공되므로 공동 파티셔닝이 필요 없어요. 즉 각 인스턴스는 체인지로그 스트림의 전체 사본을 가져요. 또한 KeyValueMapper가 KStream에서 GlobalKTable로의 키가 아닌 조인을 허용해요. KTable-KTable Foreign-Key 조인도 공동 파티셔닝을 요구하지 않아요. Kafka Streams가 내부적으로 Foreign-Key 조인에 공동 파티셔닝을 보장해요.
참고: Kafka Streams는 공동 파티셔닝 요구사항을 부분적으로만 검증해요: 파티션 할당 단계, 즉 런타임에서 조인 양쪽의 파티션 수가 같은지 Kafka Streams가 검증해요. 같지 않으면
TopologyBuilderException(런타임 예외)이 던져져요. Kafka Streams는 조인의 입력 스트림/테이블 간에 파티셔닝 전략이 일치하는지 검증할 수 없다는 점을 유의해요 — 사용자가 보장해야 해요.
데이터 공동 파티셔닝 보장: 조인의 입력이 아직 공동 파티셔닝되지 않았다면 수동으로 보장해야 해요. 아래에 설명된 것과 같은 절차를 따를 수 있어요. 병목을 피하기 위해 더 적은 파티션의 토픽을 더 큰 파티션 수에 맞게 리파티셔닝하는 것이 권장돼요. 기술적으로는 더 많은 파티션의 토픽을 더 작은 파티션 수로 리파티셔닝하는 것도 가능해요. 스트림-테이블 조인의 경우 KStream을 리파티셔닝하는 것이 권장돼요. KTable을 리파티셔닝하면 두 번째 상태 저장소가 생길 수 있기 때문이에요. 테이블-테이블 조인의 경우 KTable 크기를 고려해 더 작은 KTable을 리파티셔닝할 수도 있어요.
- 조인에서 기본 카프카 토픽이 더 적은 파티션 수를 가진 입력 KStream/KTable을 식별. 이 스트림/테이블을 "SMALLER"라고 하고 조인의 다른 쪽을 "LARGER"라고 하자. 카프카 토픽의 파티션 수를 알려면 예: CLI 도구
bin/kafka-topics와--describe옵션을 사용. - 애플리케이션에서 "SMALLER"의 데이터를 재파티셔닝.
repartition으로 재파티셔닝할 때 "LARGER"와 같은 파티셔너를 사용해야 해요.- "SMALLER"가 KStream이면:
KStream#repartition(Repartitioned.numberOfPartitions(...)). - "SMALLER"가 KTable이면:
KTable#toStream#repartition(Repartitioned.numberOfPartitions(...).toTable()).
- "SMALLER"가 KStream이면:
- 애플리케이션에서 "LARGER"와 새 스트림/테이블 사이의 조인을 수행.
KStream-KStream 조인
KStream-KStream 조인은 항상 윈도우 조인이에요. 그렇지 않으면 조인을 수행하는 데 사용되는 내부 상태 저장소(예: 슬라이딩 윈도우 또는 "버퍼")의 크기가 무한정 커지기 때문이에요.
참고 (헤더 인지 상태 저장소): dsl.store.format=HEADERS에서는 inner 스트림-스트림 조인이 헤더 인지 조인 윈도우 저장소를 사용해요. Left와 outer 스트림-스트림 조인은 아직 일치하지 않은 레코드에 대해 별도의 버퍼 저장소도 사용하며, 그 버퍼는 헤더 인지가 아니에요. 이 버퍼를 통과하는 레코드는 헤더를 잃어요. 조인 출력 레코드는 계산되거나 병합된 헤더를 가지지 않아요. 현재 전달 경로는 출력을 트리거한 레코드의 헤더를 가질 수 있어요. 스트림-스트림 조인의 경우 한 쪽의 새 입력 레코드가 다른 쪽의 각 일치 레코드에 대해 조인 출력을 생성하고, 주어진 조인 윈도우에 그러한 일치 레코드가 여러 개 있을 수 있다는 점을 강조하는 것이 중요해요.
조인 출력 레코드는 사용자 제공 ValueJoiner를 활용해 다음과 같이 효과적으로 생성돼요:
KeyValue<K, LV> leftRecord = ...;
KeyValue<K, RV> rightRecord = ...;
ValueJoiner<LV, RV, JV> joiner = ...;
KeyValue<K, JV> joinOutputRecord = KeyValue.pair(
leftRecord.key, /* by definition, leftRecord.key == rightRecord.key */
joiner.apply(leftRecord.value, rightRecord.value)
);
| 변환 | 설명 |
|---|---|
Inner Join (windowed) (KStream, KStream) -> KStream |
이 스트림을 다른 스트림과 INNER JOIN 수행. 이 연산은 윈도우지만, 조인된 스트림은 KStream<Windowed<K>, ...>가 아닌 KStream<K, ...> 타입. 데이터는 공동 파티셔닝되어야 함. |
Left Join (windowed) (KStream, KStream) -> KStream |
이 스트림을 다른 스트림과 LEFT JOIN 수행. |
Outer Join (windowed) (KStream, KStream) -> KStream |
이 스트림을 다른 스트림과 OUTER JOIN 수행. |
// Inner Join (windowed) 예시
import java.time.Duration;
KStream<String, Long> left = ...;
KStream<String, Double> right = ...;
KStream<String, String> joined = left.join(right,
(leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue, /* ValueJoiner */
JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5)),
Joined.with(
Serdes.String(), /* key */
Serdes.Long(), /* left value */
Serdes.Double()) /* right value */
);
조인은 키 기반(leftRecord.key == rightRecord.key)이고 윈도우 기반이에요. 즉 사용자 제공 JoinWindows가 정의한 대로 두 입력 레코드의 타임스탬프가 서로 "가깝다"면 조인돼요. 윈도우는 레코드 타임스탬프에 대한 추가 조인 프레디킷을 정의해요. 새 입력이 수신될 때마다 아래 조건에서 조인이 트리거돼요:
null키 또는null값을 가진 입력 레코드는 무시되고 조인을 트리거하지 않아요.
Left/outer 조인의 경우 한 쪽에서 다른 쪽과 일치하지 않는 각 입력 레코드에 대해 ValueJoiner가 ValueJoiner#apply(leftRecord.value, null)(Left) 또는 null 조합(Outer)으로 호출돼요. 이러한 결과는 지정된 유예 기간(grace period)이 지난 후 발행돼요. 주의: 폐기된 JoinWindows.of(...).grace(...) API를 사용하면 조급하게 가짜 left/right 결과가 발행될 수 있어요.
KTable-KTable Equi-Join
KTable-KTable equi-조인은 항상 비윈도우 조인이에요. 관계형 데이터베이스의 대응물과 일관되도록 설계되었어요. 두 KTable의 체인지로그 스트림은 테이블 이중성의 최신 스냅샷을 나타내기 위해 로컬 상태 저장소로 구체화돼요. 조인 결과는 조인 연산의 체인지로그 스트림을 나타내는 새 KTable이에요.
| 변환 | 설명 |
|---|---|
Inner Join (KTable, KTable) -> KTable |
이 테이블을 다른 테이블과 INNER JOIN 수행. 결과는 조인의 "현재" 결과를 나타내는 계속 업데이트되는 KTable. 데이터는 공동 파티셔닝되어야 함. |
Left Join (KTable, KTable) -> KTable |
이 테이블을 다른 테이블과 LEFT JOIN 수행. |
Outer Join (KTable, KTable) -> KTable |
이 테이블을 다른 테이블과 OUTER JOIN 수행. |
// KTable-KTable Inner Join 예시
KTable<String, Long> left = ...;
KTable<String, Double> right = ...;
KTable<String, String> joined = left.join(right,
(leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue /* ValueJoiner */
);
상세 동작:
null키를 가진 입력 레코드는 무시되고 조인을 트리거하지 않아요.null값을 가진 입력 레코드는 해당 키에 대한 톰스톤 으로 해석되며, 이는 키가 테이블에서 삭제되었음을 나타내요. 톰스톤은 조인을 트리거하지 않아요. 입력 톰스톤이 수신되면 필요할 때(즉 해당 키가 조인 결과 KTable에 이미 존재하는 경우에만) 출력 톰스톤이 조인 결과 KTable로 직접 전달돼요.- 버전 테이블을 조인할 때 순서가 뒤집힌 입력 레코드, 즉 같은 테이블에서 같은 키와 더 큰 타임스탬프를 가진 다른 레코드가 이미 처리된 경우, 해당 레코드는 무시되고 조인을 트리거하지 않아요.
KTable-KTable Foreign-Key Join
KTable-KTable 외래 키 조인은 항상 비윈도우 조인이에요. 외래 키 조인은 SQL의 조인과 유사해요. 대략적인 예:
SELECT ... FROM {this KTable} JOIN {other KTable} ON {other.key} = {result of foreignKeyExtractor(this.value)} ...
연산의 출력은 조인 결과를 포함하는 새 KTable이에요. 두 KTable의 체인지로그 스트림은 최신 스냅샷을 나타내기 위해 로컬 상태 저장소로 구체화돼요. 외래 키 추출기 함수가 왼쪽 레코드에 적용되어 새 중간 레코드가 생성되고, 오른쪽 테이블의 해당 기본 키와 조회·조인하는 데 사용돼요. 결과는 조인 연산의 체인지로그 스트림을 나타내는 새 KTable이에요.
왼쪽 KTable은 오른쪽 KTable의 같은 키에 매핑되는 여러 레코드를 가질 수 있어요. 단일 왼쪽 KTable 항목의 업데이트는 해당 키가 오른쪽 KTable에 존재하는 경우 단일 출력 이벤트를 만들 수 있어요. 결과적으로 오른쪽 KTable 항목의 단일 업데이트는 같은 외래 키를 가진 왼쪽 KTable의 각 레코드에 대한 업데이트를 만들 거예요.
// KTable-KTable Foreign-Key Inner Join 예시
KTable<String, Long> left = ...;
KTable<Long, Double> right = ...;
//This foreignKeyExtractor simply uses the left-value to map to the right-key.
Function<Long, Long> foreignKeyExtractor = (v) -> v;
//Alternative: with access to left table key
BiFunction<String, Long, Long> foreignKeyExtractor = (k, v) -> v;
KTable<String, String> joined = left.join(right, foreignKeyExtractor,
(leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue /* ValueJoiner */
);
상세 동작:
- 조인은 키 기반, 즉
foreignKeyExtractor.apply(leftRecord.value) == rightRecord.key. foreignKeyExtractor가null을 생성하는 레코드는 무시되고 조인을 트리거하지 않아요.null외래 키와 조인하려면 적절한 센티널 값(즉 String 필드의 경우"NULL", 자동 증가 정수 필드의 경우-1)을 사용해 조인해요.null값을 가진 입력 레코드는 톰스톤으로 해석. 톰스톤은 조인을 트리거하지 않아요.
KStream-KTable Join
KStream-KTable 조인은 항상 비윈도우 조인이에요. KStream(레코드 스트림)에서 새 레코드를 받으면 KTable(체인지로그 스트림)에 대한 테이블 조회를 수행할 수 있게 해줘요. 예시 사용 사례는 사용자 활동 스트림(KStream)을 최신 사용자 프로필 정보(KTable)로 강화(enrich)하는 것이에요.
| 변환 | 설명 |
|---|---|
Inner Join (KStream, KTable) -> KStream |
이 스트림을 테이블과 INNER JOIN 수행, 효과적으로 테이블 조회. 데이터는 공동 파티셔닝되어야 함. |
Left Join (KStream, KTable) -> KStream |
이 스트림을 테이블과 LEFT JOIN 수행, 효과적으로 테이블 조회. |
// KStream-KTable Inner Join 예시
KStream<String, Long> left = ...;
KTable<String, Double> right = ...;
KStream<String, String> joined = left.join(right,
(leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue, /* ValueJoiner */
Joined.keySerde(Serdes.String()) /* key */
.withValueSerde(Serdes.Long()) /* left value */
.withGracePeriod(Duration.ZERO) /* grace period */
);
상세 동작:
- 스트림(왼쪽)의 입력 레코드만 조인을 트리거해요. 테이블(오른쪽)의 입력 레코드는 내부 오른쪽 조인 상태만 업데이트해요.
- 스트림의
null키 또는null값을 가진 입력 레코드는 무시되고 조인을 트리거하지 않아요. - 테이블의
null값을 가진 입력 레코드는 톰스톤으로 해석. 톰스톤은 조인을 트리거하지 않아요.
테이블이 버전화되면, 조인할 테이블 레코드는 timestamped 조회를 수행해 결정돼요. 즉 조인되는 테이블 레코드는 스트림 레코드 타임스탬프보다 작거나 같은 타임스탬프를 가진 최신-타임스탬프 레코드가 돼요. 스트림 레코드 타임스탬프가 테이블의 히스토리 리텐션보다 오래되면 레코드는 버려져요. 유예 기간을 사용하려면 테이블이 버전화되어야 해요.
KStream-GlobalKTable Join
KStream-GlobalKTable 조인은 항상 비윈도우 조인이에요. KStream(레코드 스트림)에서 새 레코드를 받으면 GlobalKTable(전체 체인지로그 스트림)에 대한 테이블 조회를 수행할 수 있게 해줘요. 예시 사용 사례는 사용자 활동 스트림(KStream)을 최신 사용자 프로필 정보(GlobalKTable)와 추가 맥락 정보(추가 GlobalKTable)로 강화하는 "스타 쿼리" 또는 "스타 조인"이에요. 그러나 GlobalKTable은 시간의 개념이 없으므로 KStream-GlobalKTable 조인은 시간 조인(temporal join)이 아니고, GlobalKTable 업데이트와 KStream 레코드 처리 사이의 이벤트 시간 동기화가 없어요.
높은 수준에서 KStream-GlobalKTable 조인은 KStream-KTable 조인과 매우 유사해요. 그러나 글로벌 테이블은 파티셔닝된 테이블과 비교해 어느 정도 비용을 들여 훨씬 더 많은 유연성을 제공해요:
- 데이터 공동 파티셔닝을 요구하지 않아요.
- 효율적인 "스타 조인" — 대규모 "팩트" 스트림을 "차원" 테이블과 조인하는 것을 허용해요.
- 외래 키와의 조인을 허용해요. 스트림 레코드의 키뿐 아니라 레코드 값의 데이터로도 테이블 데이터를 조회할 수 있어요.
- 심하게 편향된 데이터로 작업해야 하고 따라서 핫 파티션으로 고통받는 많은 사용 사례를 가능하게 해요.
- 연속으로 여러 조인을 수행해야 할 때 파티셔닝된 KTable 대응물보다 보통 더 효율적이에요.
// KStream-GlobalKTable Inner Join 예시
KStream<String, Long> left = ...;
GlobalKTable<Integer, Double> right = ...;
KStream<String, String> joined = left.join(right,
(leftKey, leftValue) -> leftKey.length(), /* derive a (potentially) new key by which to lookup against the table */
(leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue /* ValueJoiner */
);
상세 동작:
GlobalKTable은KafkaStreams인스턴스의 (재)시작 시 완전히 부트스트랩돼요. 즉 테이블은 시작 시점에 사용 가능한 기본 토픽의 모든 데이터로 완전히 채워져요. 실제 데이터 처리는 부트스트래핑이 완료된 후에야 시작돼요.- 스트림(왼쪽)의 입력 레코드만 조인을 트리거해요. 테이블(오른쪽)의 입력 레코드는 내부 오른쪽 조인 상태만 업데이트해요.
- 스트림의
null키 또는null값을 가진 입력 레코드는 무시되고 조인을 트리거하지 않아요. - 테이블의
null값을 가진 입력 레코드는 톰스톤으로 해석. 톰스톤은 조인을 트리거하지 않아요.
조인 의미론은 KStream-KTable 조인과 다르다. 왜냐하면 시간 조인이 아니기 때문이에요. 또 다른 차이는 KStream-GlobalKTable 조인의 경우 왼쪽 입력 레코드가 테이블 조회 전에 사용자 제공 KeyValueMapper로 테이블의 키스페이스에 "매핑"된다는 점이에요.
윈도잉 (Windowing)
윈도잉을 사용하면 집계나 조인 같은 상태 저장 연산을 위해 같은 키를 가진 레코드를 소위 윈도우로 그룹화하는 방법을 제어할 수 있어요. 윈도우는 레코드 키별로 추적돼요.
관련 연산은 그룹화(grouping)이며, 후속 연산을 위해 같은 키를 가진 모든 레코드를 그룹화해 데이터가 제대로 파티셔닝("키잉")되도록 해요. 일단 그룹화되면 윈도잉을 사용해 키의 레코드를 더 하위 그룹화할 수 있어요.
예를 들어 조인 연산에서 윈도잉 상태 저장소가 정의된 윈도우 경계 내에서 지금까지 받은 모든 레코드를 저장하는 데 사용돼요. 집계 연산에서 윈도잉 상태 저장소가 윈도우별 최신 집계 결과를 저장하는 데 사용돼요. 상태 저장소의 오래된 레코드는 지정된 윈도우 리텐션 기간 후 정리돼요. Kafka Streams는 윈도우를 적어도 이 지정된 시간 동안 유지하도록 보장해요; 기본값은 하루이며 Materialized#withRetention()으로 변경할 수 있어요.
DSL이 지원하는 윈도우 유형:
| 윈도우 이름 | 동작 | 간단 설명 |
|---|---|---|
| Hopping 시간 윈도우 | 시간 기반 | 고정 크기, 겹치는 윈도우 |
| Tumbling 시간 윈도우 | 시간 기반 | 고정 크기, 겹치지 않는, 틈 없는 윈도우 |
| Sliding 시간 윈도우 | 시간 기반 | 레코드 타임스탬프 간의 차이에 동작하는 고정 크기 겹치는 윈도우 |
| 세션 윈도우 | 세션 기반 | 동적으로 크기 조정되는, 겹치지 않는, 데이터 중심 윈도우 |
Hopping 시간 윈도우
Hopping 시간 윈도우는 시간 간격 기반의 윈도우예요. 고정 크기의 (가능하면) 겹치는 윈도우를 모델링해요. Hopping 윈도우는 두 속성으로 정의돼요: 윈도우의 크기와 전진 간격(일명 "hop"). 전진 간격은 윈도우가 이전 윈도우에 비해 얼마나 앞으로 이동하는지를 지정해요. 예를 들어 크기 5분, 전진 간격 1분의 hopping 윈도우를 구성할 수 있어요. Hopping 윈도우는 겹칠 수 있고 — 일반적으로 겹치므로 — 데이터 레코드는 하나 이상의 그러한 윈도우에 속할 수 있어요.
참고: Hopping 윈도우 vs Sliding 윈도우 — Hopping 윈도우는 다른 스트림 처리 도구에서 "sliding windows"라고도 불려요. Kafka Streams는 sliding 윈도우의 의미론이 hopping 윈도우와 다른 학술 문헌의 용어를 따르고 있어요.
다음 코드는 크기 5분, 전진 간격 1분의 hopping 윈도우를 정의해요:
import java.time.Duration;
import org.apache.kafka.streams.kstream.TimeWindows;
// A hopping time window with a size of 5 minutes and an advance interval of 1 minute.
// The window's name -- the string parameter -- is used to e.g. name the backing state store.
Duration windowSize = Duration.ofMinutes(5);
Duration advance = Duration.ofMinutes(1);
TimeWindows.ofSizeWithNoGrace(windowSize).advanceBy(advance);
Hopping 시간 윈도우는 epoch에 정렬되며, 하한 경계는 포함(inclusive)이고 상한 경계는 제외(exclusive)예요. "Epoch에 정렬"은 첫 번째 윈도우가 타임스탬프 0에서 시작한다는 뜻이에요. 예를 들어 크기 5000ms, 전진 간격 3000ms의 hopping 윈도우는 예측 가능한 경계 [0;5000),[3000;8000),...를 가져요.
이전에 본 비윈도우 집계와 달리, 윈도우 집계는 키 타입이 Windowed<K>인 윈도우 KTable을 반환해요. 이것은 다른 윈도우에서 같은 키를 가진 집계 값을 구분하기 위해서예요. 해당 윈도우 인스턴스와 임베드된 키는 각각 Windowed#window()와 Windowed#key()로 검색할 수 있어요.
Tumbling 시간 윈도우
Tumbling 시간 윈도우는 hopping 시간 윈도우의 특수한 경우이며, 후자처럼 시간 간격 기반의 윈도우예요. 고정 크기, 겹치지 않는, 틈 없는 윈도우를 모델링해요. Tumbling 윈도우는 단일 속성으로 정의돼요: 윈도우의 크기. Tumbling 윈도우는 윈도우 크기가 전진 간격과 같은 hopping 윈도우예요. Tumbling 윈도우는 절대 겹치지 않으므로, 데이터 레코드는 정확히 하나의 윈도우에만 속해요.
Tumbling 시간 윈도우는 epoch에 정렬되며 하한은 포함, 상한은 제외예요. 예: 크기 5000ms의 tumbling 윈도우는 경계 [0;5000),[5000;10000),...를 가져요.
다음 코드는 크기 5분의 tumbling 윈도우를 정의해요:
import java.time.Duration;
import org.apache.kafka.streams.kstream.TimeWindows;
// A tumbling time window with a size of 5 minutes (and, by definition, an implicit
// advance interval of 5 minutes), and grace period of 1 minute.
Duration windowSize = Duration.ofMinutes(5);
Duration gracePeriod = Duration.ofMinutes(1);
TimeWindows.ofSizeAndGrace(windowSize, gracePeriod);
// The above is equivalent to the following code:
TimeWindows.ofSizeAndGrace(windowSize, gracePeriod).advanceBy(windowSize);
Sliding 시간 윈도우
Sliding 윈도우는 실제로 hopping·tumbling 윈도우와 꽤 다르다. Kafka Streams에서 sliding 윈도우는 JoinWindows 클래스로 지정되는 조인 연산과 SlidingWindows 클래스로 지정되는 윈도우 집계에 사용돼요.
Sliding 윈도우는 시간 축 위에서 지속적으로 미끄러지는 고정 크기 윈도우를 모델링해요. 이 모델에서 두 데이터 레코드는 (대칭 윈도우의 경우) 타임스탬프의 차이가 윈도우 크기 내에 있으면 같은 윈도우에 포함된다고 말해요. Sliding 윈도우가 시간 축을 따라 이동하면서 레코드가 sliding 윈도우의 여러 스냅샷에 들어갈 수 있지만, 레코드의 각 고유 조합은 하나의 sliding 윈도우 스냅샷에만 나타나요.
다음 코드는 시간 차이 10분, 유예 기간 30분의 sliding 윈도우를 정의해요:
import org.apache.kafka.streams.kstream.SlidingWindows;
// A sliding time window with a time difference of 10 minutes and grace period of 30 minutes
Duration timeDifference = Duration.ofMinutes(10);
Duration gracePeriod = Duration.ofMinutes(30);
SlidingWindows.ofTimeDifferenceAndGrace(timeDifference, gracePeriod);
Sliding 윈도우는 epoch가 아니라 데이터 레코드 타임스탬프에 정렬돼요. hopping·tumbling 윈도우와 달리 sliding 윈도우의 하한·상한 시간 경계는 모두 포함(inclusive)이에요.
세션 윈도우
세션 윈도우는 키 기반 이벤트를 소위 세션(session)으로 집계하는 데 사용되며, 그 과정을 세션화(sessionization)라고 해요. 세션은 정의된 비활동 갭(또는 "유휴")으로 구분되는 활동 기간을 나타내요. 기존 세션의 비활동 갭 내에 해당하는 처리된 모든 이벤트는 기존 세션으로 병합돼요. 이벤트가 세션 갭 밖에 해당하면 새 세션이 생성돼요.
세션 윈도우는 다른 윈도우 유형과 다음에서 다르다:
- 모든 윈도우는 키 전체에 걸쳐 독립적으로 추적돼요 — 예: 다른 키의 윈도우는 보통 다른 시작·종료 시간을 가져요.
- 윈도우 크기가 다양해요 — 같은 키의 윈도우도 보통 크기가 달라요.
세션 윈도우의 주요 적용 분야는 사용자 행동 분석이에요. 세션 기반 분석은 뉴스 웹사이트·소셜 플랫폼에서의 사용자 방문 수 같은 단순한 메트릭에서 고객 전환 퍼널·이벤트 흐름 같은 더 복잡한 메트릭까지 다양할 수 있어요.
다음 코드는 비활동 갭 5분의 세션 윈도우를 정의해요:
import java.time.Duration;
import org.apache.kafka.streams.kstream.SessionWindows;
// A session window with an inactivity gap of 5 minutes.
SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(5));
윈도우 최종 결과 (Window Final Results)
Kafka Streams에서 윈도우 계산은 결과를 지속적으로 업데이트해요. 윈도우에 새 데이터가 도착하면 새로 계산된 결과가 다운스트림으로 발행돼요. 많은 애플리케이션에서 새 결과가 항상 사용 가능하기 때문에 이상적이며, Kafka Streams는 지속적 계산 프로그래밍을 매끄럽게 만들도록 설계되었어요. 그러나 일부 애플리케이션은 윈도우 계산의 최종 결과에 대해서만 조치를 취해야 해요. 흔한 예는 알림 보내기 또는 업데이트를 지원하지 않는 시스템에 결과 전달하기예요.
예를 들어 사용자별 시간당 윈도우 카운트가 있다고 가정해보세요. 사용자가 한 시간에 3개 미만의 이벤트를 가질 때 알림을 보내고 싶다면 실제 도전이 있어요. 처음에는 모든 사용자가 이벤트를 충분히 쌓을 때까지 이 조건에 일치할 거예요. 따라서 누군가 조건에 일치할 때 단순히 알림을 보낼 수 없어요; 특정 윈도우에 대해 더 이상 이벤트를 볼 수 없다는 것을 알 때까지 기다렸다가 그런 다음 알림을 보내야 해요.
Kafka Streams는 이 로직을 정의하는 깔끔한 방법을 제공해요: 윈도우 계산을 정의한 후 중간 결과를 억제(suppress)하고, 윈도우가 닫힐 때 각 사용자의 최종 카운트를 발행할 수 있어요.
예를 들어:
KGroupedStream<UserId, Event> grouped = ...;
grouped
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofHours(1), Duration.ofMinutes(10)))
.count()
.suppress(Suppressed.untilWindowCloses(unbounded()))
.filter((windowedUserId, count) -> count < 3)
.toStream()
.foreach((windowedUserId, count) -> sendAlert(windowedUserId.window(), windowedUserId.key(), count));
이 프로그램의 핵심 부분:
ofSizeAndGrace(Duration.ofHours(1), Duration.ofMinutes(10))— 지정된 10분 유예 기간은 윈도우가 받아들일 이벤트의 지연을 제한할 수 있게 해줘요. 예: 09:00~10:00 윈도우는 10:10까지 순서가 뒤집힌 레코드를 받고, 그 시점에 윈도우가 닫혀요..suppress(Suppressed.untilWindowCloses(...))— 억제 연산자가 윈도우가 닫힐 때까지 아무것도 발행하지 않도록 구성하고, 그런 다음 최종 결과를 발행해요. 예: 사용자U가 09:00~10:10 사이에 10개 이벤트를 받으면, 억제의 다운스트림filter는 10:10까지 윈도우 키U@09:00-10:00에 대해 이벤트를 받지 않고, 그런 다음 값10을 가진 정확히 하나를 받아요. 이것이 윈도우 카운트의 최종 결과예요.unbounded()— 윈도우가 닫힐 때까지 이벤트를 저장하는 데 사용되는 버퍼를 구성해요. 프로덕션 코드는 버퍼에 사용할 메모리 양에 상한을 둘 수 있지만, 이 간단한 예는 상한이 없는 버퍼를 만들어요.
억제는 다른 Kafka Streams 연산자와 같아요. count에서 나오는 두 개의 분기(하나는 억제, 하나는 비억제, 또는 심지어 여러 다른 구성의 억제)가 있는 토폴로지를 만들 수 있어요. 필요한 곳에 억제를 적용하고 그 외에는 기본 지속 업데이트 동작에 의존할 수 있어요.
참고 (헤더 인지 상태 저장소): suppress()는 헤더 인지가 아닌 인메모리 버퍼를 사용해요. 업스트림 레코드에 첨부된 레코드 헤더는 KIP-1285에 따라 dsl.store.format=HEADERS가 전역으로 설정되어 있어도 억제 경계를 넘어 보존되지 않아요.
자세한 정보는 Suppressed 구성 객체의 JavaDoc과 KIP-328을 참고해요.
프로세서 적용 (Processor API 통합)
앞서 언급한 무상태·상태 저장 변환 외에 DSL에서 Processor API를 활용할 수도 있어요. 유용할 수 있는 시나리오:
- 커스터마이즈: DSL에 아직 없거나 아직 없는 특수·커스터마이즈된 로직을 구현해야 할 때.
- 사용 용이성과 필요할 때의 완전한 유연성 결합: 일반적으로 DSL의 표현력을 선호하지만, 처리의 특정 단계는 DSL이 제공하는 것보다 더 많은 유연성과 튜닝을 요구해요. 예를 들어 레코드의 토픽, 파티션, 오프셋 정보 같은 메타데이터에 접근하는 것은 Processor API만 가능해요. 그러나 그것 때문에 완전히 Processor API로 전환하고 싶지는 않아요.
- 다른 도구에서 마이그레이션: 명령형 API를 제공하는 다른 스트림 처리 기술에서 마이그레이션하고, 레거시 코드 일부를 Processor API로 옮기는 것이 바로 DSL로 완전히 마이그레이션하는 것보다 빠르거나 쉬운 경우.
연산과 개념:
KStream#process:Processor(주어진ProcessorSupplier가 제공)를 적용해 스트림의 모든 레코드를 한 번에 하나씩 처리.KStream#processValues:FixedKeyProcessor(주어진FixedKeyProcessorSupplier가 제공)를 적용해 스트림의 모든 레코드를 한 번에 하나씩 처리.Processor: 키-값 쌍 레코드의 프로세서.ContextualProcessor:ProcessorContext인스턴스를 관리하는Processor의 추상 구현.FixedKeyProcessor: 키가 불변인 키-값 쌍 레코드의 프로세서.ContextualFixedKeyProcessor:FixedKeyProcessorContext인스턴스를 관리하는FixedKeyProcessor의 추상 구현.ProcessorSupplier: 하나 이상의Processor인스턴스를 만들 수 있는 프로세서 공급자.FixedKeyProcessorSupplier: 하나 이상의FixedKeyProcessor인스턴스를 만들 수 있는 프로세서 공급자.
예시는 로그 심각도별 분류(process, 무상태), 텍스트 메시지의 속어 교체(processValues, 무상태), 충성도 프로그램 누적 할인(process, 상태 저장), 교통 레이더 차량 카운트(processValues, 상태 저장) 등이 포함돼요. 각각 ContextualProcessor 또는 FixedKeyProcessor를 구현하고 context().forward()로 레코드를 다운스트림으로 발행해요.
// 프로세서 적용 예시 — 로그 심각도별 분류
public class CategorizingLogsBySeverityExample {
private static final String ERROR_LOGS_TOPIC = "error-logs-topic";
private static final String INPUT_LOGS_TOPIC = "input-logs-topic";
private static final String UNKNOWN_LOGS_TOPIC = "unknown-logs-topic";
private static final String WARN_LOGS_TOPIC = "warn-logs-topic";
public static void categorizeWithProcess(final StreamsBuilder builder) {
final KStream<String, String> logStream = builder.stream(INPUT_LOGS_TOPIC);
logStream.process(LogSeverityProcessor::new)
.to((key, value, recordContext) -> {
// Determine the target topic dynamically
if ("ERROR".equals(key)) return ERROR_LOGS_TOPIC;
if ("WARN".equals(key)) return WARN_LOGS_TOPIC;
return UNKNOWN_LOGS_TOPIC;
});
}
private static class LogSeverityProcessor extends ContextualProcessor<String, String, String, String> {
@Override
public void process(final Record<String, String> record) {
if (record.value() == null) {
return; // Skip null values
}
// Assume the severity is the first word in the log message
// For example: "ERROR: Disk not found" -> "ERROR"
final int colonIndex = record.value().indexOf(':');
final String severity = colonIndex > 0 ? record.value().substring(0, colonIndex).trim() : "UNKNOWN";
// Route logs based on severity
switch (severity) {
case "ERROR":
context().forward(record.withKey(ERROR_LOGS_TOPIC));
break;
case "WARN":
context().forward(record.withKey(WARN_LOGS_TOPIC));
break;
case "INFO":
// INFO logs are ignored
break;
default:
// Forward to an "unknown" topic for logs with unrecognized severities
context().forward(record.withKey(UNKNOWN_LOGS_TOPIC));
}
}
}
}
Transformers 제거 및 프로세서로의 마이그레이션
Kafka 4.0부터 Kafka Streams API의 transform, flatTransform, transformValues, flatTransformValues, process 같은 몇 가지 폐기된 메서드가 제거됐어요. 이 메서드들은 더 다재다능한 Processor API로 대체되었어요.
다음 폐기된 메서드는 더 이상 Kafka Streams에서 사용할 수 없어요:
KStream#transformKStream#flatTransformKStream#transformValuesKStream#flatTransformValuesKStream#process
Processor API가 이제 이 모든 메서드에 대한 통합 대체물로 기능해요. 무상태와 상태 저장 연산 모두에 대한 지원을 유지하면서 API 표면을 단순화해요.
주의:
KStream.transformValues()나KStream.flatTransformValues()를 사용하고 "merge repartition topics" 최적화를 활성화한 경우, 프로그램을KStream.processValues()로 다시 작성하는 것은 KAFKA-19668 때문에 안전하지 않을 수 있어요. 이 경우 Kafka Streams 4.0.0이나 4.1.0으로 업그레이드하지 말고, 수정이 포함된 Kafka Streams 4.0.1이나 4.1.1을 사용해야 해요. 수정은 하위 호환성을 위해 기본으로 활성화되지 않으므로 구성__enable.process.processValue.fix__ = true를 설정하고StreamsBuilder()생성자에 전달해 수정을 활성화해야 해요.
final Properties properties = new Properties();
properties.put(StreamsConfig.APPLICATION_ID_CONFIG, ...);
properties.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, ...);
properties.put(TopologyConfig.InternalConfig.ENABLE_PROCESS_PROCESSVALUE_FIX, true);
final StreamsBuilder builder = new StreamsBuilder(new TopologyConfig(new StreamsConfig(properties)));
processValues()로의 재작성이 올바른지, 그리고 어떤 비호환성도 도입하지 않는지 검증하기 위해 이전·새 토폴로지의 Topology.describe() 출력을 비교하는 것이 권장돼요. 또한 비프로덕션 환경에서 업그레이드를 테스트해야 해요.
마이그레이션은 flatTransform → process, flatTransformValues → processValues, transform → process, transformValues → processValues와 같이 대응돼요. 각 경우 Processor 또는 FixedKeyProcessor 인터페이스 구현을 제공해요.
또한 '이전' Processor API를 DSL에 통합했던 process 메서드(Processor 대신 새 api.Processor)도 제거되었다는 점을 언급할 가치가 있어요. 예시로 페이지뷰가 미리 정의된 인기도 임계값(예: 1000뷰)에 도달하면 이메일 알림을 보내는 시스템이 있고, 새 process로의 마이그레이션을 보여줘요.
Controlling KTable emit rate
KTable은 논리적으로 지속적으로 업데이트되는 테이블이에요. 이 업데이트는 새 데이터가 사용 가능할 때마다 다운스트림 연산자로 전달돼, 전체 계산이 가능한 한 최신 상태가 되도록 보장해요. 논리적으로 말하면 대부분의 프로그램은 일련의 변환을 설명하며 업데이트율은 프로그램 동작에 영향을 미치지 않아요. 이러한 경우 업데이트율은 성능 문제에 가까워요. 연산자는 커밋 간격과 배치 크기 구성을 조정해 네트워크 트래픽(카프카 브로커로)과 디스크 트래픽(로컬 상태 저장소로)을 모두 최적화할 수 있어요.
그러나 일부 애플리케이션은 외부 시스템 호출 같은 다른 조치를 취해야 하므로, 예를 들어 KStream#foreach의 호출율을 제어해야 해요. KTable 레코드 캐시의 부작용으로 달성하기보다 KTable#suppress 연산자로 직접 속도 제한을 부과할 수 있어요.
예:
KGroupedTable<String, String> groupedTable = ...;
groupedTable
.count()
.suppress(untilTimeLimit(ofMinutes(5), maxBytes(1_000_000L).emitEarlyWhenFull()))
.toStream()
.foreach((key, count) -> updateCountsDatabase(key, count));
이 구성은 updateCountsDatabase가 각 key에 대해 5분에 한 번 이하로만 이벤트를 받도록 보장해요. 각 키의 최신 상태가 그 5분 동안 메모리에서 버퍼링되어야 한다는 점을 유의해요. 이 버퍼에 사용할 최대 메모리 양(이 경우 1MB)을 제어할 수 있는 옵션이 있어요. 레코드 수로 제한을 부과하는 옵션도 있어요 (또는 두 제한 모두 지정하지 않을 수 있음). 추가로 버퍼가 가득 차면 어떻게 할지 선택하는 것이 가능해요. 이 예는 완화된 접근을 취해 5분 제한 전에 가장 오래된 레코드를 발행해 버퍼를 다시 크기로 되돌려요. 대안으로 처리를 중지하고 애플리케이션을 종료하도록 선택할 수 있어요. 이것은 극단적으로 보일 수 있지만, 5분 제한이 절대적으로 강제될 것이라는 보장을 제공해요. 애플리케이션 종료 후 버퍼에 더 많은 메모리를 할당하고 처리를 재개할 수 있어요. 조기 발행(emit early)이 대부분의 애플리케이션에서 선호돼요.
자세한 정보는 Suppressed 구성 객체의 JavaDoc과 KIP-328을 참고해요.
테이블 프로세서에 타임스탬프 기반 의미론 사용
기본적으로 Kafka Streams의 테이블은 오프셋 기반 의미론(offset-based semantics)을 사용해요. 같은 키에 대해 여러 레코드가 도착하면 가장 큰 레코드 오프셋을 가진 것이 키의 최신 레코드로 간주되고, 이것이 테이블에 대해 계산된 집계·조인 결과에 나타나는 레코드예요. 이것은 순서가 뒤집힌 데이터가 있어도 마찬가지예요. 가장 큰 오프셋을 가진 레코드가 키의 최신 레코드로 간주돼요, 그 레코드가 가장 큰 타임스탬프를 가지지 않아도요.
오프셋 기반 의미론의 대안은 타임스탬프 기반 의미론이에요. 타임스탬프 기반 의미론에서 가장 큰 타임스탬프를 가진 레코드가 최신 레코드로 간주돼요, 더 큰 오프셋(그리고 더 작은 타임스탬프)을 가진 다른 레코드가 있어도요. (키별) 순서가 뒤집힌 데이터가 없으면 오프셋 기반 의미론과 타임스탬프 기반 의미론은 동등해요; 차이는 순서가 뒤집힌 데이터가 있을 때만 나타나요.
Kafka Streams 3.5부터 버전 상태 저장소를 통해 타임스탬프 기반 의미론을 지원해요. 테이블이 버전 상태 저장소로 구체화되면 버전 테이블이고 순서가 뒤집힌 데이터가 있을 때 다른 프로세서 의미론을 가져요.
- 스트림-테이블 조인을 수행할 때 스트림 측 레코드는 스트림 레코드 타임스탬프보다 작거나 같은 타임스탬프를 가진 최신-타임스탬프 테이블 레코드와 조인돼요.
- 테이블에 대해 계산된 집계는 최신-오프셋 레코드 대신 각 키의 최신-타임스탬프 레코드를 포함해요. (키별) 순서가 뒤집힌 업데이트는 새 집계 결과를 트리거하지 않아요. 이것은
count·reduce연산과aggregate연산에서 모두 참이에요. - 테이블 조인은 최신-오프셋 레코드 대신 각 키의 최신-타임스탬프 레코드를 사용해요. (키별) 순서가 뒤집힌 업데이트는 새 조인 결과를 트리거하지 않아요. 이것은 기본 키 테이블-테이블 조인과 외래 키 테이블-테이블 조인 모두에서 참이에요.
- 테이블 필터 연산은 더 이상 연속 톰스톤을 억제하지 않아요. 따라서 비버전 테이블을 필터링할 때보다 다운스트림에서 더 많은
null레코드를 관찰할 수 있어요. 이것은 순서가 뒤집힌 데이터의 경우 완전한 버전 이력을 다운스트림에 보존하기 위해서예요. suppress연산은 버전 테이블에서 허용되지 않아요. 버전 이력을 붕괴시켜 정의되지 않은 동작을 초래할 수 있기 때문이에요.
테이블이 버전 저장소로 구체화되면 다음 중 하나가 발생할 때까지 다운스트림 테이블도 버전화된 것으로 간주돼요:
- 다운스트림 테이블이 (비버전 저장소 공급자로 또는 저장소 공급자 없이) 명시적으로 구체화됨 (모든 저장소는 기본적으로 비버전이며 기본 저장소 공급자 포함).
- 집계·조인을 포함한 어떤 상태 저장 변환이 발생.
- 테이블이 스트림으로 변환되고 다시 되돌아옴.
일부 프로세서의 결과는 버전 저장소로 구체화되어서는 안 되는데, 그 프로세서들이 완전한 이전 버전 이력을 생성하지 않아 버전 테이블로의 구체화가 예측할 수 없는 결과를 초래할 수 있기 때문이에요:
- 테이블·스트림 집계 모두에 대한 집계 프로세서.
aggregate,count,reduce연산 포함. - 기본 키·외래 키 조인을 모두 포함한 테이블-테이블 조인 프로세서.
버전 저장소에 대한 자세한 내용과 애플리케이션에서 사용을 시작하는 방법은 여기를 참고해요.
스트림을 카프카로 다시 쓰기
모든 스트림과 테이블은 카프카 토픽으로 (지속적으로) 다시 쓸 수 있어요. 출력 데이터는 상황에 따라 카프카로 가는 도중에 재파티셔닝될 수 있어요.
| 카프카로 쓰기 | 설명 |
|---|---|
To KStream -> void |
종단 연산. 레코드를 카프카 토픽에 씀. Produced 클래스를 통해 Serdes를 명시적으로 지정해야 할 수 있음. Produced 인스턴스로 데이터가 어떻게 생성되는지 지정하는 변형(예: 출력 레코드가 출력 토픽의 파티션에 분산되는 방식을 제어하는 StreamPartitioner), TopicNameExtractor 인스턴스로 각 레코드에 대해 보낼 토픽을 동적으로 선택하는 변형도 있음. |
KStream<String, Long> stream = ...;
// Write the stream to the output topic, using the configured default key
// and value serdes.
stream.to("my-stream-output-topic");
// Write the stream to the output topic, using explicit key and value serdes,
// (thus overriding the defaults in the config properties).
stream.to("my-stream-output-topic", Produced.with(Serdes.String(), Serdes.Long()));
다음 조건 중 하나가 참이면 데이터 재파티셔닝이 발생해요:
- 출력 토픽이 스트림/테이블과 다른 파티션 수를 가질 때.
KStream이 재파티셔닝으로 표시되었을 때.- 출력 레코드를 출력 토픽의 파티션에 분산하는 것을 명시적으로 제어하는 커스텀
StreamPartitioner를 제공할 때. - 출력 레코드의 키가
null일 때.
참고: 카프카가 아닌 시스템에 쓰고 싶을 때 — 데이터를 카프카로 다시 쓰는 것 외에도 처리 끝에 커스텀 프로세서를 스트림 싱크로 적용해 예를 들어 외부 데이터베이스에 쓸 수 있어요. 먼저, 그렇게 하는 것은 권장되는 패턴이 아니에요 — Kafka Connect API를 사용하는 것을 강력히 권장해요. 그러나 그러한 싱크 프로세서를 사용한다면, 그러한 외부 시스템과 통신할 때 메시지 전달 의미론을 보장하는 것(전달 실패 시 재시도, 메시지 중복 방지 등)이 이제 여러분의 책임이라는 점을 인지하세요.
Streams 애플리케이션 테스트
Kafka Streams는 애플리케이션 테스트를 돕는 test-utils 모듈을 함께 제공해요.
Kafka Streams DSL for Scala
⚠️ 폐기 안내: Kafka Streams DSL for Scala 라이브러리(
kafka-streams-scala)는 Kafka 4.3부터 폐기되고 Kafka 5.0에서 제거될 예정이에요. 마이그레이션 가이드와 Scala 래퍼에서 Java API로 마이그레이션하는 방법을 보여주는 코드 예시를 참고해요.
Kafka Streams DSL Java API는 빌더 디자인 패턴을 기반으로 해요. 이 API는 Scala에서 호출될 수 있지만 몇 가지 문제가 있어요:
- 추가 타입 주석: Java API는 Scala 컴파일러의 타입 추론과 완전히 호환되지 않는 방식으로 Java 제네릭을 사용해요. 따라서 사용자가 Scala 코드에 타입 주석을 추가해야 하는데, 이는 Scala에서 다소 관용적이지 않아 보여요.
- 장황함(Verbosity): 일부 경우 Java API는 관용적 Scala에 비해 너무 장황해 보여요.
- 타입 안전성 부족: Java API는 컴파일 시간 타입 안전성이 때때로 무너지고 런타임 오류가 발생할 수 있는 일부 옵션을 제공해요. 구성의 일부로 정의된 Serdes가 컴파일 시간에 타입 체크되지 않는다는 사실에서 비롯돼요. 따라서 누락된 Serdes는 런타임 오류를 초래할 수 있어요.
Kafka Streams DSL for Scala 라이브러리는 위에서 제기된 우려를 해결하는 기존 Java API를 위한 래퍼예요. 처음부터 개발된 Scala 라이브러리에서 구현할 관용적 Scala API를 제공하려고 시도하지 않아요. 의도는 더 나은 타입 추론, 향상된 표현력, 더 적은 보일러플레이트를 통해 Scala에서 Java API를 더 유용하게 만드는 것이에요.
이 라이브러리는 Scala Stream DSL Java API를 래핑해 다음을 제공해요:
- Scala에서 더 나은 타입 추론.
- 애플리케이션 코드에서 더 적은 보일러플레이트.
- 개발자가 원래 Java API로 얻는 일반적인 빌더 스타일 구성.
- 더 나은 추상화와 덜 장황함을 이끄는 암묵적(implicit) 시리얼라이저·디시리얼라이저.
- 컴파일 시간에 더 나은 타입 안전성.
Kafka Streams DSL for Scala가 제공하는 모든 기능은 루트 패키지 이름 org.apache.kafka.streams.scala 아래에 있어요. Java API의 많은 공개 타입이 래핑되어 있으며, 다음 Scala 추상화가 사용자에게 제공돼요:
org.apache.kafka.streams.scala.StreamsBuilderorg.apache.kafka.streams.scala.kstream.KStreamorg.apache.kafka.streams.scala.kstream.KTableorg.apache.kafka.streams.scala.kstream.KGroupedStreamorg.apache.kafka.streams.scala.kstream.KGroupedTableorg.apache.kafka.streams.scala.kstream.SessionWindowedKStreamorg.apache.kafka.streams.scala.kstream.TimeWindowedKStream
라이브러리는 또한 올바른 의미론을 위해 사용자가 사용해야 하는 여러 유틸리티 추상화·모듈이 있어요.
org.apache.kafka.streams.scala.ImplicitConversions: Scala와 Java 클래스 사이의 암묵적 변환을 범위로 가져오는 모듈.org.apache.kafka.streams.scala.serialization.Serdes: 암묵적으로 가져올 수 있는 모든 기본 Serdes와 커스텀 Serdes를 만들기 위한 헬퍼를 포함하는 모듈.
라이브러리는 Scala 2.12와 2.13으로 크로스 빌드돼요. Scala 2.13으로 컴파일된 라이브러리를 참조하려면 maven pom.xml에 다음을 포함해요:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams-scala_2.13</artifactId>
<version>4.3.1</version>
</dependency>
Scala 2.12으로 컴파일된 라이브러리를 사용하려면 artifactId를 kafka-streams-scala_2.12로 바꿔요. SBT를 사용할 때는 다음으로 올바른 라이브러리를 참조할 수 있어요:
libraryDependencies += "org.apache.kafka" %% "kafka-streams-scala" % "4.3.1"
샘플 사용법
라이브러리는 카프카 스트림즈의 원래 Java 추상화를 Scala 래퍼 객체 안에 래핑하고 그들 사이에 암묵적 변환을 사용함으로써 동작해요. 모든 Scala 추상화는 해당 Java 추상화와 동일하게 이름이 지정되지만 라이브러리의 다른 패키지에 있어요. 예: Scala 클래스 org.apache.kafka.streams.scala.StreamsBuilder는 org.apache.kafka.streams.StreamsBuilder의 래퍼이고, org.apache.kafka.streams.scala.kstream.KStream은 org.apache.kafka.streams.kstream.KStream의 래퍼이죠.
다음은 Java KStream의 래퍼인 KStream 인스턴스를 만드는 Scala StreamsBuilder를 사용하는 고전적인 WordCount 프로그램 예시예요. 그런 다음 테이블로 구체화하고 Java KTable의 래퍼인 다시 KTable을 얻어요.
import java.time.Duration
import java.util.Properties
import org.apache.kafka.streams.kstream.Materialized
import org.apache.kafka.streams.scala.ImplicitConversions._
import org.apache.kafka.streams.scala._
import org.apache.kafka.streams.scala.kstream._
import org.apache.kafka.streams.{KafkaStreams, StreamsConfig}
object WordCountApplication extends App {
import Serdes._
val props: Properties = {
val p = new Properties()
p.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application")
p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker1:9092")
p
}
val builder: StreamsBuilder = new StreamsBuilder
val textLines: KStream[String, String] = builder.stream[String, String]("TextLinesTopic")
val wordCounts: KTable[String, Long] = textLines
.flatMapValues(textLine => textLine.toLowerCase.split("\\W+"))
.groupBy((_, word) => word)
.count(Materialized.as("counts-store"))
wordCounts.toStream.to("WordsWithCountsTopic")
val streams: KafkaStreams = new KafkaStreams(builder.build(), props)
streams.start()
sys.ShutdownHookThread {
streams.close(Duration.ofSeconds(10))
}
}
위 코드 스니펫에서 어떤 Serdes, Grouped, Produced, Consumed 또는 Joined도 명시적으로 제공하지 않아요. 그것들은 또한 구성에 지정된 어떤 Serdes에도 의존하지 않아요. 사실 구성에 지정된 모든 Serdes는 Scala API에 의해 무시돼요. 모든 Serdes와 Grouped, Produced, Consumed 또는 Joined는 아래 Implicit Serdes 섹션에서 논의되는 암묵적 Serdes를 통해 처리돼요. 구성 기반 Serdes로부터의 완전한 독립성이 이 라이브러리를 완전히 타입 안전하게 만드는 것이에요. Serdes, Grouped, Produced, Consumed 또는 Joined의 누락된 인스턴스는 컴파일 시간 오류로 표시돼요.
Implicit Serdes
Java API에 대한 Scala 사용자의 흔한 불만 중 하나는 API 호출에서 Serdes의 반복적인 사용이에요. 많은 API가 Grouped, Produced, Repartitioned, Consumed 또는 Joined 같은 추상화를 통해 Serdes를 받아야 해요. 그리고 사용자는 이 클래스들의 with 함수를 통해 매번 공급해야 해요.
이 라이브러리는 Scala 암묵적 파라미터의 힘을 사용해 이 우려를 완화해요. 사용자는 암묵적 Serdes 또는 Grouped, Produced, Repartitioned, Consumed 또는 Joined의 암묵적 값을 한 번 제공할 수 있고 코드를 덜 장황하게 만들 수 있어요. 사실 암묵적 Serdes만 범위에 있으면 라이브러리가 Grouped, Produced, Consumed 또는 Joined의 인스턴스를 범위에서 사용 가능하게 만들어요.
라이브러리는 또한 일반적으로 사용되는 기본 타입의 모든 암묵적 Serdes를 Scala 모듈에 번들해요 — 모듈 vals를 가져오기만 하면 모든 Serdes가 범위에 있어요. 모듈식 암묵값의 유사한 전략이 사용자 정의 Serdes에도 채택될 수 있어요 (작성자 정의 Serdes는 다음 섹션에서 논의).
예:
// DefaultSerdes brings into scope implicit Serdes (mostly for primitives)
// that will set up all Grouped, Produced, Consumed and Joined instances.
// So all APIs below that accept Grouped, Produced, Consumed or Joined will
// get these instances automatically
import Serdes._
val builder = new StreamsBuilder()
val userClicksStream: KStream[String, Long] = builder.stream(userClicksTopic)
val userRegionsTable: KTable[String, String] = builder.table(userRegionsTopic)
// The following code fragment does not have a single instance of Grouped,
// Produced, Consumed or Joined supplied explicitly.
// All of them are taken care of by the implicit Serdes imported by DefaultSerdes
val clicksPerRegion: KTable[String, Long] =
userClicksStream
.leftJoin(userRegionsTable)((clicks, region) => (if (region == null) "UNKNOWN" else region, clicks))
.map((_, regionWithClicks) => regionWithClicks)
.groupByKey
.reduce(_ + _)
clicksPerRegion.toStream.to(outputTopic)
위 코드 스니펫에서 꽤 많은 것들이 진행되고 있어요:
- 코드 스니펫은 구성에 정의된 어떤 Serdes에도 의존하지 않아요. 사실 구성의 일부로 정의된 어떤 Serdes도 무시돼요.
- 모든 Serdes는 범위의 암묵값에서 가져와요. 그리고
import Serdes._가 필요한 모든 Serdes를 범위로 가져와요. - 이것은 Java API에는 없는 컴파일 시간 타입 안전성의 예시예요.
- 코드는 덜 장황하고 데이터 스트림에서 수행하는 실제 변환에 더 집중해 보여요.
사용자 정의 Serdes
기본 기본 Serdes로 충분하지 않고 커스텀 Serdes를 정의해야 할 때 사용법은 위와 정확히 같아요. 암묵적 Serdes를 정의하고 스트림 변환 구축을 시작하기만 하면 돼요. AvroSerde가 있는 예시:
// domain object as a case class
case class UserClicks(clicks: Long)
// An implicit Serde implementation for the values we want to
// serialize as avro
implicit val userClicksSerde: Serde[UserClicks] = new AvroSerde
// Primitive Serdes
import Serdes._
// And then business as usual ..
val userClicksStream: KStream[String, UserClicks] = builder.stream(userClicksTopic)
val userRegionsTable: KTable[String, String] = builder.table(userRegionsTopic)
// Compute the total per region by summing the individual click counts per region.
val clicksPerRegion: KTable[String, Long] =
userClicksStream
// Join the stream against the table.
.leftJoin(userRegionsTable)((clicks, region) => (if (region == null) "UNKNOWN" else region, clicks.clicks))
// Change the stream from <user> -> <region, clicks> to <region> -> <clicks>
.map((_, regionWithClicks) => regionWithClicks)
// Compute the total per region by summing the individual click counts per region.
.groupByKey
.reduce(_ + _)
// Write the (continuously updating) results to the output topic.
clicksPerRegion.toStream.to(outputTopic)
사용자 정의 Serdes의 완전한 예시는 라이브러리 내부의 테스트 클래스에서 찾을 수 있어요.
더 알아보기
- Streams 개발자 가이드 — 전체 개발 문서를 봐요.
- Processor API — 저수준 API를 봐요.
- 핵심 개념 — 스트림·테이블, 시간, 윈도잉 개념을 이해해요.