
Kafka Streams 윈도우 도입기
두줄요약
재고 정산의 스파이크성 트래픽을 Kafka Streams 텀블링 윈도우로 집계했습니다. 생성 시각 추출과 파티션별 더미 이벤트로 시간 정합성과 윈도우 종료를 해결했습니다.
문제 상황
- 5분 단위 재고 변동 이력 발행과 특정 시간대 스파이크성 트래픽으로 인한 DB 부하 우려
- 생성 시각과 Kafka 발행 시각 불일치에 따른 일자별 정산 데이터 오분류 가능성
- 신규 이벤트 부재 시 스트림 시간 정지로 윈도우 결과 미발행
구조와 흐름
- 상품 코드별 5분 텀블링 윈도우와 수량 reduce 집계
- 유예 시간 내 지연 이벤트 수용 후 윈도우 종료 시 최종 집계만 발행
- 비스파이크성 데이터의 집계 없는 정산 테이블 즉시 반영
해결 방법
- TimestampExtractor로 레코드 발행 시각 대신 데이터 생성 시각을 이벤트 시간으로 추출
- Processor의 wall-clock 스케줄러에서 파티션별 외부 더미 이벤트 발행으로 스트림 시간 전진
- suppress를 통한 중간 집계 미발행과 윈도우 종료 후 단일 결과 발행
주의할 점
- 내부 forward는 소스 토픽 레코드가 아니므로 스트림 시간 갱신 불가
- WindowStore 직접 스캔·강제 발행 시 내부 윈도우 상태 미종료에 따른 중복 발행 위험
- 자정 등 시간 경계 지연 이벤트와 데이터 재계산 정책의 사전 설계 필요
