이 도메인의 학습 목표

이 도메인이 묻는 것

참고 챕터는 10장 Kafka Streams와 ksqlDB입니다. Streams는 표준 Kafka 프로듀서·컨슈머 위에 올라간 라이브러리이므로, Application Development 도메인의 지식이 그대로 적용됩니다. Streams 애플리케이션의 병렬성은 결국 입력 토픽의 파티션 수이고, 상태 저장소는 결국 Kafka의 컴팩션 토픽(changelog)입니다. 이 두 연결을 잡고 있으면 문제 대부분이 풀립니다.

핵심 개념 (1) KStream · KTable · GlobalKTable

세 추상의 의미론 차이
관점 KStream KTable GlobalKTable
해석 레코드 스트림 (append-only) changelog 스트림 (upsert) changelog 스트림 (전체 복제)
같은 키가 또 오면 별개의 이벤트 이전 값을 갱신 이전 값을 갱신
value가 null이면 값이 null인 이벤트 삭제 (tombstone) 삭제
파티션 배분 인스턴스별로 일부 파티션 인스턴스별로 일부 파티션 모든 인스턴스가 전체 복제
적합한 데이터 클릭, 결제 트랜잭션, 센서 측정, 로그 사용자 프로필, 재고 수량, 계정 잔액 국가 코드, 환율, 카테고리 마스터
크기 제약 없음 없음 (파티션 단위로 분산) 인스턴스 메모리·디스크에 전체가 들어가야 함

핵심 개념 (2) stateless · stateful · 리파티션

이 절의 표가 이 도메인 최다 출제 지점입니다. 연산을 세 축으로 분류해야 합니다 — 상태를 갖는가, 리파티션을 유발하는가, 그리고 키를 바꾸는가.

DSL 연산 분류 — 리파티션 유발 여부가 실무·시험에서 가장 중요합니다
연산 상태 키를 바꾸는가 리파티션 유발
filter · filterNotstateless아니오없음
mapValues · flatMapValuesstateless아니오없음
map · flatMapstateless가능표시됨
selectKeystateless표시됨
peek · foreachstateless아니오없음
split / branch · mergestateless아니오없음
repartitionstateless아니오즉시 발생
groupByKey아니오없음 (키가 그대로)
groupBy항상 발생
count · reduce · aggregatestateful아니오(선행 group 단계에서 결정)
모든 윈도우 연산stateful아니오
모든 조인stateful아니오앞에서 키를 바꿨으면 발생
suppressstateful아니오없음
toTablestateful아니오없음

상태 저장소와 changelog

핵심 개념 (3) 조인 매트릭스와 co-partitioning

조인 4종 — 윈도우 필요 여부와 co-partitioning 요구
조인 윈도우 co-partitioning 지원 유형 트리거
KStream ⋈ KStream 필수 (JoinWindows) 필요 inner · left · outer 양쪽 어느 레코드가 와도
KStream ⋈ KTable 불필요 필요 inner · left KStream 쪽 레코드만
KStream ⋈ GlobalKTable 불필요 불필요 inner · left KStream 쪽 레코드만
KTable ⋈ KTable 불필요 필요
(FK 조인은 불필요)
inner · left · outer 양쪽 어느 갱신이 와도

co-partitioning 요건 — 공식 요건은 2개입니다

공식 문서(streams/developer-guide/dsl-api)가 co-partitioning 요구사항으로 열거하는 것은 두 개입니다. 키가 같아야 한다는 것은 요건이 아니라 equi-join 자체의 전제로 별도 서술됩니다.

KStream-KStream 조인 — 윈도우가 필수입니다
KStream<String, Order>   orders   = builder.stream("orders");
KStream<String, Payment> payments = builder.stream("payments");

// 5분 윈도우 안에서 같은 키를 만나면 조인한다
KStream<String, Settled> settled = orders.join(
    payments,
    (order, payment) -> new Settled(order, payment),   // ValueJoiner
    JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5)),
    Joined.with(Serdes.String(), orderSerde, paymentSerde));

핵심 개념 (4) 윈도우 4종

윈도우 타입 비교
윈도우 경계 한 레코드가 속하는 윈도우 수 생성 API
Tumbling 고정 크기, 겹치지 않음 1개 TimeWindows.ofSizeWithNoGrace(size)
Hopping 고정 크기, advance 간격으로 겹침 size ÷ advance 개 TimeWindows.ofSizeAndGrace(size, grace)
.advanceBy(advance)
Sliding 레코드를 중심으로 실제 데이터가 있는 구간만 데이터에 따라 가변 SlidingWindows
.ofTimeDifferenceAndGrace(diff, grace)
Session 비활성 간격(gap)으로 구분. 크기 가변 1개 (세션 병합 가능) SessionWindows
.ofInactivityGapWithNoGrace(gap)

시간 개념과 grace period

타임스탬프 추출기의 기본값은 FailOnInvalidTimestamp이며, 음수(무효) 타임스탬프를 만나면 예외를 던집니다. 레거시 데이터를 다룰 때는 LogAndSkipOnInvalidTimestampWallclockTimestampExtractor로 바꾸는 선택지를 씁니다.

grace period는 윈도우가 닫힌 뒤에도 늦게 온 레코드를 받아 주는 유예 기간입니다. …WithNoGrace(…)를 쓰면 유예가 0이므로 지연 도착 레코드가 버려집니다. grace를 크게 주면 정확도가 오르지만 결과 확정이 늦어지고 상태가 커집니다. 중간 결과를 억제하고 최종 결과만 내보내려면 suppress()를 씁니다.

핵심 개념 (5) 병렬성과 exactly-once

인터랙티브 쿼리와 ksqlDB

상태 저장소는 KafkaStreams.store(...)로 직접 조회할 수 있습니다. 인스턴스가 여러 대면 키가 어느 인스턴스에 있는지 메타데이터로 찾아 그 인스턴스에 요청을 전달해야 합니다 (내장 RPC는 없고 애플리케이션이 구현합니다).

ksqlDB는 Kafka Streams 위에 올라간 SQL 계층이며, Apache Kafka 본체가 아니라 Confluent 컴포넌트입니다. CREATE STREAM / CREATE TABLE이 각각 KStream / KTable에 대응합니다. Java 코드 없이 스트림 처리를 하고 싶을 때 쓰고, 복잡한 로직·커스텀 상태가 필요하면 Streams로 갑니다.

반드시 외워야 할 설정값

Kafka Streams 필수 암기 설정 (Apache Kafka 4.3)
설정기본값시험 포인트
application.id(필수)컨슈머 그룹 ID · 내부 토픽 접두어 · client-id 접두어 세 역할
processing.guaranteeat_least_once값은 at_least_once / exactly_once_v2 두 개뿐
num.stream.threads1기본이 1이므로 명시하지 않으면 단일 스레드
num.standby.replicas0기본이 0이므로 복구 시 changelog 전량 재생
replication.factor-1-1이면 브로커 기본값 사용. 내부 토픽 RF를 3으로 두려면 명시
commit.interval.ms30000EOS(exactly_once_v2)에서는 더 작은 값이 적용됩니다
statestore.cache.max.bytes1048576010MiB. 0으로 두면 모든 갱신이 하위로 흘러 결과가 많아짐
state.dir${java.io.tmpdir}기본값이 임시 디렉터리 — 재시작·컨테이너 교체 시 상태 유실
default.key.serde · default.value.serdenull기본값이 없으므로 지정하지 않으면 런타임 오류
default.timestamp.extractorFailOnInvalidTimestamp무효 타임스탬프에 예외를 던짐
deserialization.exception.handlerLogAndFailExceptionHandler역직렬화 실패 시 애플리케이션이 죽습니다. 건너뛰려면 LogAndContinue…
topology.optimizationnoneall로 켜면 불필요한 리파티션 토픽을 줄일 수 있음
max.task.idle.ms0조인에서 한쪽 입력을 기다리는 시간. 0이면 기다리지 않음
buffered.records.per.partition1000파티션별 내부 버퍼
poll.ms100내부 컨슈머 poll 블로킹 시간
acceptable.recovery.lag10000이 이하로 따라잡으면 "복구 완료"로 보고 태스크를 넘김
probing.rebalance.interval.ms60000010분. warmup 진행 확인 주기

자주 나오는 함정

함정 1 — map vs mapValues

이 도메인 최다 출제 쌍
관점mapmapValues
무엇인가키와 값을 모두 바꿀 수 있음값만 바꿈
상태statelessstateless
언제 발동키가 바뀔 수 있으므로 "리파티션 필요"로 표시표시하지 않음
혼동 시 결과불필요한 리파티션 토픽 생성 → 네트워크·디스크·지연 증가
출제 형태"값만 변환하는데 리파티션을 피하려면 어떤 연산을?" → mapValues

함정 2 — KStream vs KTable

같은 토픽을 어떻게 읽느냐로 결과가 완전히 달라집니다
관점KStreamKTable
무엇인가사실의 연속키별 최신 상태
같은 키 반복각각 별개 이벤트갱신 (이전 값 무효)
언제 발동builder.stream(...)builder.table(...)
혼동 시 결과상태를 스트림으로 읽어 중복 집계이벤트를 테이블로 읽어 이벤트 유실
출제 형태"결제 이벤트 합계가 실제보다 작다" → KTable로 읽은 경우

함정 3 — KTable vs GlobalKTable

둘 다 테이블이지만 분산 방식이 다릅니다
관점KTableGlobalKTable
데이터 분산파티션 단위로 인스턴스에 분산모든 인스턴스가 전체 복제
조인 시 co-partitioning필요불필요
키가 아닌 값으로 조인불가가능 (KeyValueMapper)
혼동 시 결과큰 테이블을 GlobalKTable로 만들면 인스턴스마다 전체를 적재해 메모리·시작 시간이 폭증
출제 형태"파티션 수가 다른 참조 데이터와 조인하려면?" → GlobalKTable

함정 4 — 윈도우 4종

"겹치는가"와 "크기가 고정인가"로 구분합니다
관점TumblingHoppingSlidingSession
겹침없음있음있음없음
크기 고정고정고정고정 간격가변
경계 결정시각시각레코드 시각 기준 구간비활성 gap
혼동 시 결과Hopping을 Tumbling으로 착각하면 집계가 중복 계산됩니다
출제 형태"advance < size인 윈도우의 이름은?" → hopping

함정 5 — co-partitioning vs repartition()

하나는 요구사항, 하나는 그것을 만족시키는 수단입니다
관점co-partitioningrepartition()
무엇인가조인의 전제 조건 (키·파티션 수·파티셔너 일치)키 기준으로 다시 분배하는 연산
어디 설정인가토픽 설계 · 프로듀서 파티셔너토폴로지 코드
언제 발동조인 시 검사 (파티션 수만)호출 즉시 리파티션 토픽 생성
혼동 시 결과파티션 수가 다르면 TopologyException,
파티셔너가 다르면 조용히 누락
남용하면 불필요한 토픽·지연 증가
출제 형태"파티션 3개 토픽과 6개 토픽을 조인하려면?" → 적은 쪽을 6으로 리파티션

함정 6 — SMT vs Streams (재확인)

Connect 도메인의 같은 함정이 Streams 쪽에서도 나옵니다. 집계·조인·윈도우는 SMT로 불가능하고, 필드 하나 바꾸는 일에 Streams 애플리케이션을 띄우는 것은 과잉입니다. 경계는 "여러 레코드의 관계가 필요한가"입니다.

코드 읽기 문제 대비

스니펫 1 — 리파티션 토픽이 몇 개 생기는가

Topology.java
builder.stream("clicks")                       // key = sessionId
       .filter((k, v) -> v.getStatus() == 200)
       .mapValues(v -> v.getUrl())              // (A)
       .selectKey((k, url) -> url)              // (B)
       .groupByKey()
       .count()
       .toStream()
       .to("url-counts");

스니펫 2 — 이 조인은 동작하는가

orders 토픽은 파티션 12개, customers 토픽은 파티션 3개
KStream<String, Order>    orders    = builder.stream("orders");
KTable<String, Customer>  customers = builder.table("customers");

KStream<String, Enriched> enriched = orders.join(
    customers,
    (order, customer) -> new Enriched(order, customer));

스니펫 3 — 이 설정으로 EOS가 동작하는가

streams.properties (브로커 1대인 개발 환경)
application.id=fraud-detector
bootstrap.servers=localhost:9092
processing.guarantee=exactly_once_v2
default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
default.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde

도메인 미니 퀴즈

공식 문서 출처