CCDAK 도메인 · 가중치 12%
Kafka Streams
Streams는 별도 클러스터가 아니라 애플리케이션에 넣는 라이브러리입니다. 출제는 API 전체를 묻지 않고 세 가지 판단에 집중됩니다 — 이 데이터는 KStream인가 KTable인가, 이 연산은 상태를 갖는가 그리고 리파티션을 유발하는가, 이 조인은 co-partitioning이 필요한가.
이 도메인의 학습 목표
- KStream · KTable · GlobalKTable의 의미론 차이를 데이터 성격으로 판단할 수 있습니다.
- stateless / stateful 연산을 분류하고, 리파티션을 유발하는 연산을 골라낼 수 있습니다.
- 조인 4종의 윈도우 필요 여부와 co-partitioning 요구사항을 말할 수 있습니다.
- 윈도우 4종의 차이와 grace period의 역할을 설명할 수 있습니다.
- 태스크·스레드 병렬성과
exactly_once_v2의 전제를 알 수 있습니다.
이 도메인이 묻는 것
참고 챕터는 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 · 리파티션
이 절의 표가 이 도메인 최다 출제 지점입니다. 연산을 세 축으로 분류해야 합니다 — 상태를 갖는가, 리파티션을 유발하는가, 그리고 키를 바꾸는가.
| 연산 | 상태 | 키를 바꾸는가 | 리파티션 유발 |
|---|---|---|---|
filter · filterNot | stateless | 아니오 | 없음 |
mapValues · flatMapValues | stateless | 아니오 | 없음 |
map · flatMap | stateless | 가능 | 표시됨 |
selectKey | stateless | 예 | 표시됨 |
peek · foreach | stateless | 아니오 | 없음 |
split / branch · merge | stateless | 아니오 | 없음 |
repartition | stateless | 아니오 | 즉시 발생 |
groupByKey | — | 아니오 | 없음 (키가 그대로) |
groupBy | — | 예 | 항상 발생 |
count · reduce · aggregate | stateful | 아니오 | (선행 group 단계에서 결정) |
| 모든 윈도우 연산 | stateful | 아니오 | — |
| 모든 조인 | stateful | 아니오 | 앞에서 키를 바꿨으면 발생 |
suppress | stateful | 아니오 | 없음 |
toTable | stateful | 아니오 | 없음 |
상태 저장소와 changelog
- stateful 연산은 로컬 상태 저장소(기본 RocksDB)를 만듭니다.
- 같은 내용이 changelog 토픽(
{application.id}-{store}-changelog)에 기록됩니다. 이 토픽은cleanup.policy=compact입니다. - 인스턴스가 죽으면 다른 인스턴스가 changelog를 재생해 상태를 복구합니다.
상태가 크면 복구가 오래 걸리므로
num.standby.replicas(기본 0)를 올려 따뜻한 예비 사본을 유지합니다. - 리파티션 토픽은
{application.id}-{name}-repartition이며cleanup.policy=delete입니다 (중간 데이터이므로).
핵심 개념 (3) 조인 매트릭스와 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 자체의 전제로 별도 서술됩니다.
- 전제 — 양쪽 입력의 키가 같은 의미여야 합니다
(equi-join이 키 기준으로 수행되므로
leftKey == rightKey). - 요건 ① — 양쪽 입력 토픽의 파티션 수가 같아야 합니다. Streams가 런타임에 검증하는 유일한 항목입니다.
- 요건 ② — 양쪽에 쓰는 모든 애플리케이션이
같은 파티셔닝 전략을 써야 합니다
(같은
partitioner.class또는 같은StreamPartitioner). 검증되지 않으므로 어긋나면 예외 없이 조인이 조용히 누락됩니다.
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) |
| Sliding | 레코드를 중심으로 실제 데이터가 있는 구간만 | 데이터에 따라 가변 | SlidingWindows |
| Session | 비활성 간격(gap)으로 구분. 크기 가변 | 1개 (세션 병합 가능) | SessionWindows |
시간 개념과 grace period
- event time — 레코드에 기록된 타임스탬프. Streams의 기본 기준입니다.
- ingestion time — 브로커가 로그에 쓴 시각 (토픽
message.timestamp.type=LogAppendTime). - processing time — 처리하는 순간의 벽시계 시각.
타임스탬프 추출기의 기본값은 FailOnInvalidTimestamp이며,
음수(무효) 타임스탬프를 만나면 예외를 던집니다.
레거시 데이터를 다룰 때는 LogAndSkipOnInvalidTimestamp나
WallclockTimestampExtractor로 바꾸는 선택지를 씁니다.
grace period는 윈도우가 닫힌 뒤에도 늦게 온 레코드를 받아 주는 유예 기간입니다.
…WithNoGrace(…)를 쓰면 유예가 0이므로 지연 도착 레코드가 버려집니다.
grace를 크게 주면 정확도가 오르지만 결과 확정이 늦어지고 상태가 커집니다.
중간 결과를 억제하고 최종 결과만 내보내려면 suppress()를 씁니다.
핵심 개념 (5) 병렬성과 exactly-once
- 태스크 수 = 입력 토픽의 파티션 수입니다 (sub-topology별로 계산). 파티션이 6개면 태스크는 6개이고, 그 이상으로는 병렬 처리가 늘지 않습니다.
num.stream.threads(기본 1)는 한 인스턴스가 돌리는 스레드 수입니다. 태스크는 스레드에 배분됩니다.- 인스턴스를 늘리면 태스크가 인스턴스 사이로 재배분됩니다. 인스턴스 × 스레드 > 태스크 수가 되면 남는 스레드는 유휴 상태입니다.
- Streams는 내부적으로 컨슈머 그룹을 쓰며, 그룹 ID는
application.id입니다. 그래서application.id를 바꾸면 처음부터 다시 읽습니다.
인터랙티브 쿼리와 ksqlDB
상태 저장소는 KafkaStreams.store(...)로 직접 조회할 수 있습니다.
인스턴스가 여러 대면 키가 어느 인스턴스에 있는지 메타데이터로 찾아 그 인스턴스에 요청을 전달해야 합니다
(내장 RPC는 없고 애플리케이션이 구현합니다).
ksqlDB는 Kafka Streams 위에 올라간 SQL 계층이며,
Apache Kafka 본체가 아니라 Confluent 컴포넌트입니다.
CREATE STREAM / CREATE TABLE이 각각 KStream / KTable에 대응합니다.
Java 코드 없이 스트림 처리를 하고 싶을 때 쓰고,
복잡한 로직·커스텀 상태가 필요하면 Streams로 갑니다.
반드시 외워야 할 설정값
| 설정 | 기본값 | 시험 포인트 |
|---|---|---|
application.id | (필수) | 컨슈머 그룹 ID · 내부 토픽 접두어 · client-id 접두어 세 역할 |
processing.guarantee | at_least_once | 값은 at_least_once / exactly_once_v2 두 개뿐 |
num.stream.threads | 1 | 기본이 1이므로 명시하지 않으면 단일 스레드 |
num.standby.replicas | 0 | 기본이 0이므로 복구 시 changelog 전량 재생 |
replication.factor | -1 | -1이면 브로커 기본값 사용. 내부 토픽 RF를 3으로 두려면 명시 |
commit.interval.ms | 30000 | EOS(exactly_once_v2)에서는 더 작은 값이 적용됩니다 |
statestore.cache.max.bytes | 10485760 | 10MiB. 0으로 두면 모든 갱신이 하위로 흘러 결과가 많아짐 |
state.dir | ${java.io.tmpdir} | 기본값이 임시 디렉터리 — 재시작·컨테이너 교체 시 상태 유실 |
default.key.serde · default.value.serde | null | 기본값이 없으므로 지정하지 않으면 런타임 오류 |
default.timestamp.extractor | FailOnInvalidTimestamp | 무효 타임스탬프에 예외를 던짐 |
deserialization.exception.handler | LogAndFailExceptionHandler | 역직렬화 실패 시 애플리케이션이 죽습니다. 건너뛰려면 LogAndContinue… |
topology.optimization | none | all로 켜면 불필요한 리파티션 토픽을 줄일 수 있음 |
max.task.idle.ms | 0 | 조인에서 한쪽 입력을 기다리는 시간. 0이면 기다리지 않음 |
buffered.records.per.partition | 1000 | 파티션별 내부 버퍼 |
poll.ms | 100 | 내부 컨슈머 poll 블로킹 시간 |
acceptable.recovery.lag | 10000 | 이 이하로 따라잡으면 "복구 완료"로 보고 태스크를 넘김 |
probing.rebalance.interval.ms | 600000 | 10분. warmup 진행 확인 주기 |
자주 나오는 함정
함정 1 — map vs mapValues
| 관점 | map | mapValues |
|---|---|---|
| 무엇인가 | 키와 값을 모두 바꿀 수 있음 | 값만 바꿈 |
| 상태 | stateless | stateless |
| 언제 발동 | 키가 바뀔 수 있으므로 "리파티션 필요"로 표시 | 표시하지 않음 |
| 혼동 시 결과 | 불필요한 리파티션 토픽 생성 → 네트워크·디스크·지연 증가 | — |
| 출제 형태 | "값만 변환하는데 리파티션을 피하려면 어떤 연산을?" → mapValues | |
함정 2 — KStream vs KTable
| 관점 | KStream | KTable |
|---|---|---|
| 무엇인가 | 사실의 연속 | 키별 최신 상태 |
| 같은 키 반복 | 각각 별개 이벤트 | 갱신 (이전 값 무효) |
| 언제 발동 | builder.stream(...) | builder.table(...) |
| 혼동 시 결과 | 상태를 스트림으로 읽어 중복 집계 | 이벤트를 테이블로 읽어 이벤트 유실 |
| 출제 형태 | "결제 이벤트 합계가 실제보다 작다" → KTable로 읽은 경우 | |
함정 3 — KTable vs GlobalKTable
| 관점 | KTable | GlobalKTable |
|---|---|---|
| 데이터 분산 | 파티션 단위로 인스턴스에 분산 | 모든 인스턴스가 전체 복제 |
| 조인 시 co-partitioning | 필요 | 불필요 |
| 키가 아닌 값으로 조인 | 불가 | 가능 (KeyValueMapper) |
| 혼동 시 결과 | 큰 테이블을 GlobalKTable로 만들면 인스턴스마다 전체를 적재해 메모리·시작 시간이 폭증 | |
| 출제 형태 | "파티션 수가 다른 참조 데이터와 조인하려면?" → GlobalKTable | |
함정 4 — 윈도우 4종
| 관점 | Tumbling | Hopping | Sliding | Session |
|---|---|---|---|---|
| 겹침 | 없음 | 있음 | 있음 | 없음 |
| 크기 고정 | 고정 | 고정 | 고정 간격 | 가변 |
| 경계 결정 | 시각 | 시각 | 레코드 시각 기준 구간 | 비활성 gap |
| 혼동 시 결과 | Hopping을 Tumbling으로 착각하면 집계가 중복 계산됩니다 | |||
| 출제 형태 | "advance < size인 윈도우의 이름은?" → hopping | |||
함정 5 — co-partitioning vs repartition()
| 관점 | co-partitioning | repartition() |
|---|---|---|
| 무엇인가 | 조인의 전제 조건 (키·파티션 수·파티셔너 일치) | 키 기준으로 다시 분배하는 연산 |
| 어디 설정인가 | 토픽 설계 · 프로듀서 파티셔너 | 토폴로지 코드 |
| 언제 발동 | 조인 시 검사 (파티션 수만) | 호출 즉시 리파티션 토픽 생성 |
| 혼동 시 결과 | 파티션 수가 다르면 TopologyException,파티셔너가 다르면 조용히 누락 | 남용하면 불필요한 토픽·지연 증가 |
| 출제 형태 | "파티션 3개 토픽과 6개 토픽을 조인하려면?" → 적은 쪽을 6으로 리파티션 | |
함정 6 — SMT vs Streams (재확인)
Connect 도메인의 같은 함정이 Streams 쪽에서도 나옵니다. 집계·조인·윈도우는 SMT로 불가능하고, 필드 하나 바꾸는 일에 Streams 애플리케이션을 띄우는 것은 과잉입니다. 경계는 "여러 레코드의 관계가 필요한가"입니다.
코드 읽기 문제 대비
스니펫 1 — 리파티션 토픽이 몇 개 생기는가
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 — 이 조인은 동작하는가
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가 동작하는가
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
이어서 볼 곳
도메인 미니 퀴즈
공식 문서 출처
- Streams Core Concepts — 스트림·테이블 이원성, 태스크와 스레드, 시간 개념
- Streams DSL — 연산 분류, 리파티션, 조인 매트릭스, co-partitioning 요구사항과
TopologyException, 윈도우 4종 API - Streams Configs —
num.stream.threads,num.standby.replicas,state.dir,processing.guarantee기본값 - Memory Management —
statestore.cache.max.bytes - Interactive Queries — 상태 저장소 직접 조회
- Streams Configs (표) — 전체 설정 기본값