
Kafka Streams는 카프카 클러스터에 저장된 데이터를 입력/출력으로 삼아 “애플리케이션 내부에서” 스트림 처리를 수행하는 라이브러리라는 점이 핵심입니다. Source
Kafka Streams는 “Kafka 토픽을 입력으로 받아(consume) 변환/집계/조인/윈도우 처리 등을 수행한 뒤, 결과를 다시 Kafka 토픽으로 내보내(produce) 실시간 파이프라인을 만드는 JVM 라이브러리”입니다. Source
별도 클러스터(예: 별도의 스트림 엔진)를 운영하기보다, 일반 서비스처럼 배포/스케일링하면서 스트림 처리를 붙일 수 있습니다. Source
Kafka Streams는 로컬 상태 저장소를 활용해 윈도우 집계, 조인 같은 stateful 연산을 효율적으로 처리할 수 있습니다. Source
실패/재시도 상황에서도 레코드가 중복 처리되지 않도록 exactly-once processing semantics를 지원합니다. Source
DSL 문서 기준으로 Kafka Streams는 KStream/KTable/GlobalKTable 같은 추상화를 제공합니다. Source
여러 연산(map/filter/join/aggregate)을 연결한 처리 그래프가 토폴로지입니다.
(공식 개념에서도 stateful 연산의 핵심으로 윈도우 조인/집계를 언급) Source
“최근 5분”, “10초 단위”처럼 시간 구간으로 묶어 집계할 때 사용합니다.
아래는 “Kafka Streams로 딱 맞게” 해결되는 대표 케이스들입니다.
상황: 결제 이벤트에서 FRAUD_SUSPECT만 따로 토픽으로 보내고, 나머지는 정상 처리로 분기.
왜 Streams?: 단순하지만 초고속/항상 켜진 파이프라인이 필요할 때 적합.
관련 개념: DSL 기반의 스트림/테이블 추상화 Source
상황: 클릭 이벤트를 상품ID로 그룹화해 1분 단위로 카운트 집계.
포인트: 이건 stateful(상태 필요)이고, 보통 windowing을 곁들입니다.
Kafka Streams는 윈도우 집계 같은 stateful 연산을 빠르게 지원한다고 명시합니다. Source
상황: “10분 내 결제 5회 이상 + 국가 변경” 같은 룰로 경보 이벤트 생성
구현 힌트: KTable(최신 상태) + windowed aggregation 조합이 흔합니다.
상황: 주문 이벤트와 결제 이벤트를 일정 시간 창(window) 안에서 매칭해 “주문완료” 이벤트 생성.
왜 Streams?: 스트림-스트림 조인은 전형적인 stateful + window 기반 패턴입니다. Source
상황: 여러 이벤트(회원수정/주소변경/탈퇴)를 소비해 회원 프로필의 최신 상태 뷰를 토픽/스토어로 유지
효과: 읽기 API는 “현재 상태”를 빠르게 제공하고, 쓰기는 이벤트로만 기록(이벤트 소싱 스타일)
(위 영상 링크는 검색 결과에서 확인된 video_id 기반입니다.)