Kafka Streams: 토픽 안의 스트림을 처리하는 라이브러리
Kafka Streams: 토픽 안의 스트림을 처리하는 라이브러리
Kafka Streams는 Kafka에 저장된 이벤트 스트림을 실시간으로 처리하고 변환하는 클라이언트 라이브러리예요. 별도 클러스터가 필요 없이, 기존 Kafka 토픽을 읽고 쓰는 애플리케이션 형태로 스트림 처리 애플리케이션과 마이크로서비스를 만들 수 있어요.
스트림 처리 애플리케이션 만들기
Kafka Streams는 프로듀서·컨슈머 위에 한 단계 높은 함수형 DSL을 제공해요. 이벤트를 토픽에서 읽어 변환·집계·조인하고, 결과를 다시 토픽에 쓰는 흐름을 표현할 수 있어요. 별도 처리 클러스터를 운영할 필요가 없다는 게 가장 큰 특징이에요.
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> lines = builder.stream("streams-plaintext-input");
lines.flatMapValues(value -> Arrays.asList(value.toLowerCase(Locale.getDefault()).split("\\W+")))
.groupBy((key, word) -> word)
.count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store"))
.toStream()
.to("streams-wordcount-output");
KStream과 KTable: 스트림을 보는 두 눈
Kafka Streams에는 두 가지 추상화가 있어요. KStream은 이벤트(레코드)의 흐름 그 자체로, 새 이벤트가 오면 변환·분기·병합하고 싶을 때 써요. KTable은 스트림의 상태(테이블 관점)로, 같은 키의 최신 값만 유지하는 업데이트 로그처럼 동작해요. 집계·조인의 기준으로 삼기 좋아요.
KTable<String, Long> counts = builder.table("user-activity");
상태 저장과 윈도우 처리
Kafka Streams는 상태 저장(stateful) 처리를 기본으로 지원해요. 집계·카운트·조인은 내부 상태 스토어에 중간 결과를 보관하고, 필요하면 로컬 RocksDB를 백업 저장소로 써요. 또 **윈도우(window)**를 지정해 이벤트 시간(Event Time) 기준으로 일정 구간 동안의 집계를 수행할 수 있어요. 예를 들어 5분 단위 팀블링 윈도우로 페이지 조회수를 세는 식이에요.
더 알아보기
- 스트림 처리를 시작하는 실습은 Streams Quick Start 문서를 보세요.
- DSL의 핵심 개념은 Streams Core Concepts 문서를 보세요.
- 토픽과 파티션 기본기는 Quickstart 문서를 보세요.