Kafka Quick Start (2) — Kafka Streams 처리, 환경 종료, Ecosystem
Apache Kafka Official Documentation — Getting Started (v4.3) — Apache Software Foundation Getting Started 공식문서 — 3. Quick Start Step 7~8, 4. Ecosystem (pp.13-15)
이 모듈을 다 읽으면
- Kafka Streams 라이브러리가 제공하는 기능(상태 유지 연산, 윈도잉, 조인, exactly-once)을 나열할 수 있다
- 공식 WordCount 예제 코드에서 stream → flatMapValues → groupBy → count → toStream/to로 이어지는 각 단계가 하는 일을 설명할 수 있다
- 로컬 Kafka 환경을 안전하게 종료하고, 필요하면 로컬 데이터까지 삭제하는 절차를 설명할 수 있다
Quick Start의 마지막 단계인 Kafka Streams를 이용한 이벤트 처리(WordCount 예제)와 로컬 환경 종료 절차, 그리고 아주 짧게 소개되는 4번째 섹션 Ecosystem을 함께 정리한다. Kafka Streams는 Producer/Consumer API 위에 얹힌 고수준 Java/Scala 클라이언트 라이브러리로, 서버 쪽 클러스터 기술의 이점(확장성·탄력성·내결함성·분산)을 표준 애플리케이션 코드 작성의 단순함과 결합한다.
Kafka Streams로 이벤트 처리하기 — WordCount 예제
데이터가 Kafka에 이벤트로 저장되고 나면, Java/Scala용 클라이언트 라이브러리인 Kafka Streams로 그 데이터를 처리할 수 있다. 입력·출력 데이터가 Kafka 토픽에 저장되는 미션 크리티컬한 실시간 애플리케이션과 마이크로서비스를 구현할 수 있게 해주는데, 클라이언트 쪽에서 표준 Java·Scala 애플리케이션을 작성·배포하는 단순함과, Kafka의 서버 쪽 클러스터 기술이 주는 이점(고도의 확장성·탄력성·내결함성·분산)을 결합한 것이 핵심이다. 이 라이브러리는 정확히 한 번(exactly-once) 처리, 상태 유지(stateful) 연산과 집계, 윈도잉, 조인, 이벤트-타임 기반 처리 등을 지원한다.
대표적인 첫 예제로 널리 알려진 WordCount 알고리즘을 Kafka Streams로 구현한 코드는 다음과 같다.
KStream<String, String> textLines = builder.stream("quickstart-events");
KTable<String, Long> wordCounts = textLines
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split(" ")))
.groupBy((keyIgnored, word) -> word)
.count();
wordCounts.toStream().to("output-topic", Produced.with(Serdes.String(), Serdes.Long()));
코드를 단계별로 보면, builder.stream("quickstart-events")는 quickstart-events 토픽을 읽어 KStream<String, String> textLines를 만든다 — 각 레코드가 하나의 텍스트 줄에 해당하는, 끝없이 이어지는 이벤트 스트림이다. 이어서 flatMapValues로 각 줄의 값을 소문자로 바꾸고 공백 기준으로 단어들로 쪼개, 레코드 하나(한 줄)를 여러 개의 단어 레코드로 펼친다(flat-map). groupBy((keyIgnored, word) -> word)는 원래 키는 무시하고 단어 자체를 새 키로 삼아 레코드를 재그룹화하며, count()는 그룹별(=단어별)로 몇 번 등장했는지 세어 KTable<String, Long> wordCounts를 만든다 — KTable은 각 단어 키에 대한 '현재까지의 누적 카운트'를 나타내는, 계속 갱신되는 테이블이다. 마지막으로 wordCounts.toStream()으로 이 테이블의 변경 이력을 다시 이벤트 스트림으로 바꾸고, .to("output-topic", Produced.with(Serdes.String(), Serdes.Long()))로 output-topic이라는 토픽에 (단어, 누적 카운트) 쌍을 문자열 키·Long 값의 Serde로 직렬화해 게시한다.
핵심 포인트
- Kafka Streams는 Producer/Consumer 위의 고수준 Java/Scala 라이브러리로, 클라이언트 코드의 단순함 + 서버 클러스터의 확장성/탄력성/내결함성을 결합한다
- 지원 기능: exactly-once 처리, 상태 유지 연산·집계, 윈도잉, 조인, 이벤트-타임 처리
- WordCount 흐름: builder.stream(토픽)으로 KStream 생성 → flatMapValues로 줄을 소문자화·공백 분리해 단어 레코드로 펼침 → groupBy(단어)로 재그룹화 → count()로 KTable 생성
- wordCounts.toStream().to("output-topic", Produced.with(Serdes.String(), Serdes.Long()))로 결과 KTable을 변경 스트림으로 바꿔 다시 토픽에 게시한다
Kafka 환경 종료하기
Quick Start를 마쳤다면 환경을 정리해도 되고, 계속 실습을 이어가도 된다. 정리하려면 먼저 아직 켜 둔 프로듀서·컨슈머 클라이언트를 Ctrl-C로 멈추고, 이어서 Kafka 브로커도 Ctrl-C로 멈춘다. 만약 지금까지 만든 이벤트를 포함해 로컬 Kafka 환경의 데이터까지 모두 지우고 싶다면 다음 명령을 실행한다.
$ rm -rf /tmp/kafka-logs /tmp/kraft-combined-logs
핵심 포인트
- 종료 순서: ① 프로듀서·컨슈머 클라이언트를 Ctrl-C로 정지 → ② Kafka 브로커를 Ctrl-C로 정지
- 로컬 데이터까지 완전히 지우려면 rm -rf /tmp/kafka-logs /tmp/kraft-combined-logs 실행 (KRaft 모드의 결합 로그 디렉토리 포함)
Ecosystem — 메인 배포판 밖의 통합 도구들
공식 문서의 Ecosystem 섹션은 매우 짧다. Kafka 메인 배포판 바깥에서 Kafka와 통합되는 다양한 도구가 많이 있으며, ecosystem 페이지가 스트림 처리 시스템, Hadoop 연동, 모니터링, 배포 도구 등 이런 도구 다수를 목록으로 정리해 안내한다는 내용만 담겨 있다.
핵심 포인트
- Kafka 메인 배포판 밖에도 스트림 처리·Hadoop 연동·모니터링·배포 등을 위한 서드파티 통합 도구가 다수 존재하며, 공식 ecosystem 페이지가 이를 목록으로 안내한다