빠른참조 · Streams
Streams 치트시트
DSL 연산자 표의 마지막 열 "리파티션 유발"이 이 페이지의 핵심입니다.
리파티션은 내부 토픽을 하나 더 만들고 네트워크 왕복을 한 번 더 만들기 때문에,
실무에서 Streams 성능 문제의 대부분이 여기서 생깁니다.
map과 mapValues의 차이가 곧 성능 차이입니다.
3초 요약
- 키를 바꿀 수 있는 연산은 리파티션 대상으로 마킹됩니다 —
map,flatMap,selectKey,groupBy. - 마킹된 스트림에 그룹화나 조인이 오면 그때 실제로 리파티션 토픽이 생깁니다.
- 키를 안 바꾸는
mapValues/flatMapValues/filter는 마킹되지 않습니다. - 조인은 co-partitioning이 전제입니다. 예외는 GlobalKTable 조인과 KTable FK 조인뿐입니다.
DSL 연산자 표
| 연산 | 입력 → 출력 | 상태 저장 | 리파티션 유발 | 메모 |
|---|---|---|---|---|
| 소스 · 싱크 | ||||
stream() |
입력 토픽 → KStream |
무상태 | 안 함 | 레코드 스트림으로 읽습니다 (append-only 의미론) |
table() |
입력 토픽 → KTable |
상태 | 안 함 | changelog 의미론. 같은 키의 새 값이 이전 값을 덮어씁니다 |
globalTable() |
입력 토픽 → GlobalKTable |
상태 | 안 함 | 모든 인스턴스가 전체 파티션의 사본을 갖습니다. 그래서 co-partitioning이 불필요합니다 |
to() |
KStream → void |
— | 조건부 | 출력 토픽 파티션 수가 다르거나, 스트림이 마킹됐거나, 커스텀 StreamPartitioner를 주거나, 키가 null이면 리파티션합니다 |
| 무상태 변환 (stateless) | ||||
filter() |
KStream→KStream, KTable→KTable |
무상태 | 안 함 | 조건이 참인 레코드만 남깁니다 |
filterNot() |
KStream→KStream, KTable→KTable |
무상태 | 안 함 | 조건이 거짓인 레코드만 남깁니다 |
map() |
KStream→KStream |
무상태 | 마킹함 | 키를 바꿀 수 있으므로 리파티션 대상이 됩니다. 값만 바꾸면 mapValues를 쓰세요 |
mapValues() |
KStream→KStream, KTable→KTable |
무상태 | 안 함 | 키를 보존하므로 리파티션이 필요 없습니다. map보다 항상 우선 |
flatMap() |
KStream→KStream |
무상태 | 마킹함 | 0..N개 레코드로 확장. 키 변경 가능 |
flatMapValues() |
KStream→KStream |
무상태 | 안 함 | 키를 보존한 채 값을 확장. flatMap보다 우선 |
selectKey() |
KStream→KStream |
무상태 | 마킹함 | 키를 명시적으로 바꿉니다. 조인 전 co-partitioning을 맞출 때 씁니다 |
split() / Branch |
KStream→BranchedKStream |
무상태 | 안 함 | 술어를 순서대로 평가해 첫 매치 하나의 스트림으로 보냅니다. 매치가 없으면 default 브랜치 또는 폐기 |
merge() |
KStream→KStream |
무상태 | 안 함 | 두 스트림을 합칩니다. 순서는 보장되지 않습니다 |
peek() |
KStream→KStream |
무상태 | 안 함 | 부수효과를 만들면서 스트림을 그대로 통과시킵니다 (로깅·메트릭용) |
foreach() |
KStream→void, KTable→void |
무상태 | 안 함 | 종료 연산입니다. 이후 처리를 이어갈 수 없습니다 — 이어가려면 peek |
print() |
KStream→void |
무상태 | 안 함 | 디버깅용 종료 연산 |
toStream() |
KTable→KStream |
무상태 | 안 함 | 테이블을 changelog 스트림으로 봅니다 |
toTable() |
KStream→KTable |
상태 | 문서에 명시 없음 | 스트림을 테이블로 변환합니다. 앞선 연산이 키를 바꿨다면 그 마킹이 여기까지 이어집니다 |
repartition() |
KStream→KStream |
무상태 | 항상 유발 | 파티션 수를 명시적으로 지정할 수 있습니다. 자동 리파티션이 걸리지 않는 process() 계열 앞에서 필요합니다 |
| 그룹화 (집계 전 단계) | ||||
groupByKey() |
KStream→KGroupedStream |
무상태 | 마킹된 경우에만 | groupBy보다 항상 우선. 이미 원하는 키라면 리파티션이 전혀 일어나지 않습니다 |
groupBy() |
KStream→KGroupedStream, KTable→KGroupedTable |
무상태 | 항상 유발 | 키를 새로 계산하므로 무조건 리파티션합니다. 가능하면 groupByKey로 대체하세요 |
cogroup() |
KGroupedStream→CogroupedKStream |
— | 안 함 | 이미 그룹화된 스트림을 받으므로 리파티션이 그 단계에서 끝나 있습니다 |
| 집계 (상태 저장) | ||||
count() |
KGroupedStream→KTable, KGroupedTable→KTable |
상태 | 그룹화 단계에서 결정 | 키별 레코드 수 |
reduce() |
KGroupedStream→KTable, KGroupedTable→KTable |
상태 | 그룹화 단계에서 결정 | 결과 타입이 입력 값 타입과 같아야 합니다 |
aggregate() |
KGroupedStream→KTable, KGroupedTable→KTable |
상태 | 그룹화 단계에서 결정 | 결과 타입을 바꿀 수 있습니다. initializer가 필요합니다 |
windowedBy() |
KGroupedStream→TimeWindowedKStream / SessionWindowedKStream |
상태 | 그룹화 단계에서 결정 | 윈도 스토어를 씁니다. 기본 보관 기간은 1일이며 Materialized#withRetention()으로 변경합니다 |
suppress() |
KTable→KTable |
상태 | 안 함 | 윈도의 최종 결과만 내보냅니다. 중간 갱신을 억제해 downstream 부하를 줄입니다 |
| 조인 (상태 저장) | ||||
join() KStream-KStream |
(KStream, KStream)→KStream |
상태 (윈도 스토어) | 양쪽 중 마킹된 쪽을 리파티션 | 항상 윈도 조인. JoinWindows 필수. co-partitioning 필요 |
join() KStream-KTable |
(KStream, KTable)→KStream |
상태 | KStream이 마킹된 경우 | 비윈도. 스트림 레코드가 도착할 때 테이블을 조회합니다. co-partitioning 필요 |
join() KStream-GlobalKTable |
(KStream, GlobalKTable)→KStream |
상태 (전체 사본) | 안 함 | 비윈도. co-partitioning 불필요. KeyValueMapper로 키가 아닌 값으로도 조회 가능 |
join() KTable-KTable (equi) |
(KTable, KTable)→KTable |
상태 | — | 비윈도. co-partitioning 필요 |
join() KTable-KTable (FK) |
(KTable, KTable)→KTable |
상태 | Streams가 내부적으로 처리 | 비윈도. co-partitioning 불필요 — Streams가 내부에서 보장합니다. SQL의 FK 조인과 유사 |
leftJoin() |
위 5종 모두 지원 | 상태 | 각 조인과 동일 | 왼쪽 레코드는 오른쪽이 없어도 null과 함께 출력됩니다 |
outerJoin() |
KStream-KStream, KTable-KTable(equi) 만 | 상태 | 각 조인과 동일 | KStream-KTable · GlobalKTable · FK 조인에는 outer가 없습니다 |
| Processor API 혼합 | ||||
process() |
KStream→KStream |
스토어 연결 시 상태 | 자동 리파티션 안 함 | 앞에서 키를 바꿨다면 repartition()을 직접 호출해야 합니다. 이것이 process()의 가장 큰 함정 |
processValues() |
KStream→KStream |
스토어 연결 시 상태 | 안 함 | 키를 바꿀 수 없는(FixedKeyProcessor) 버전이라 리파티션이 필요 없습니다 |
map이 키를 바꿀 수 있다고 판단해 스트림이 마킹되고, groupByKey에서 리파티션 토픽이 생깁니다. 값만 바꾸는데도 네트워크 왕복이 추가됩니다.
builder.stream("orders", Consumed.with(Serdes.String(), orderSerde))
// 키를 그대로 돌려주더라도 Streams 는 알 수 없습니다.
// map 이 호출된 것만으로 "키가 바뀔 수 있다"고 마킹됩니다.
.map((k, v) -> KeyValue.pair(k, v.withTax()))
.groupByKey() // ← 여기서 리파티션 토픽 생성
.count();
mapValues는 키를 바꿀 수 없으므로 마킹되지 않고, groupByKey가 리파티션 없이 진행됩니다.
builder.stream("orders", Consumed.with(Serdes.String(), orderSerde))
// 키를 건드릴 수 없는 시그니처이므로 마킹되지 않습니다.
.mapValues(v -> v.withTax())
.groupByKey() // ← 리파티션 없음
.count();
// 1) 토폴로지를 문자열로 출력해 repartition 노드를 찾습니다
Topology topology = builder.build();
System.out.println(topology.describe());
// "Sink: KSTREAM-FILTER-...-repartition-filter" 처럼
// -repartition 이 들어간 노드가 있으면 리파티션이 생긴 것입니다
# 내부 토픽은 -- 규칙을 따릅니다
# (공식 문서는 이 규칙이 향후 릴리스에서 보장되지 않는다고 명시합니다)
bin/kafka-topics.sh --bootstrap-server localhost:9092 --list \
| grep -E "^my-streams-app-" | sort
# 출력 예
# my-streams-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
# my-streams-app-order-count-repartition
윈도 타입 4종
| 타입 | 정의 파라미터 | 겹침 | 한 레코드가 속하는 윈도 수 | 클래스 · 특징 |
|---|---|---|---|---|
| Tumbling | 크기(size)만 | 없음 | 정확히 1개 | TimeWindows. hopping의 특수형(advance = size). epoch에 정렬되어 5000ms면 경계가 [0,5000), [5000,10000)로 예측 가능합니다 |
| Hopping | 크기 + advance interval(hop) | 있음 | 여러 개 | TimeWindows. advance < size면 겹칩니다. 다른 스트림 처리 도구에서 "sliding"이라 부르는 것이 이것입니다 |
| Sliding | time difference + grace | 있음 | 여러 개 | 집계는 SlidingWindows, 조인은 JoinWindows. 두 레코드의 타임스탬프 차이가 윈도 크기 이내면 같은 윈도. 각 레코드 조합은 한 스냅샷에만 나타납니다 |
| Session | inactivity gap | 없음 (병합됨) | 1개 (병합될 수 있음) | SessionWindows. 키마다 독립적으로 추적되고 크기가 가변입니다. gap 안의 이벤트는 기존 세션에 병합되며, 늦게 온 레코드가 두 세션을 하나로 합칠 수도 있습니다 |
import java.time.Duration;
import org.apache.kafka.streams.kstream.TimeWindows;
import org.apache.kafka.streams.kstream.SlidingWindows;
import org.apache.kafka.streams.kstream.SessionWindows;
import org.apache.kafka.streams.kstream.JoinWindows;
// Tumbling — 크기 5분. advance 가 size 와 같은 hopping 이라고 볼 수 있습니다.
TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5));
// Hopping — 크기 5분, 1분마다 전진. 한 레코드가 최대 5개 윈도에 들어갑니다.
TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5))
.advanceBy(Duration.ofMinutes(1));
// Sliding (집계) — 타임스탬프 차이 10분, grace 30분
SlidingWindows.ofTimeDifferenceAndGrace(
Duration.ofMinutes(10), Duration.ofMinutes(30));
// Sliding (조인) — KStream-KStream 조인은 JoinWindows 를 씁니다
JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(10));
// Session — 비활성 간격 5분
SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(5));
조인 매트릭스
| 조인 | 윈도 | co-partitioning | INNER | LEFT | OUTER | 특징 |
|---|---|---|---|---|---|---|
| KStream ⋈ KStream | 항상 필요 (JoinWindows) |
필요 | O | O | O | 양쪽 레코드를 윈도 스토어에 보관합니다. 상태가 가장 많이 듭니다 |
| KStream ⋈ KTable | 비윈도 | 필요 | O | O | X | 스트림 레코드가 도착할 때만 조인이 트리거됩니다. 테이블 갱신은 트리거하지 않습니다 |
| KStream ⋈ GlobalKTable | 비윈도 | 불필요 | O | O | X | 모든 인스턴스가 전체 사본을 갖습니다. KeyValueMapper로 값 기반 조회가 가능해 star join에 적합. 시간 동기화가 없습니다 |
| KTable ⋈ KTable (equi) | 비윈도 | 필요 | O | O | O | 양쪽 갱신이 모두 조인을 트리거합니다 |
| KTable ⋈ KTable (foreign-key) | 비윈도 | 불필요 | O | O | X | Streams가 내부적으로 co-partitioning을 보장합니다. 왼쪽의 여러 레코드가 오른쪽의 한 키에 매핑될 수 있습니다 |
co-partitioning 요건은 2개입니다
공식 문서가 명시하는 요건은 둘입니다.
"키가 같아야 한다"는 요건이 아니라 equi-join의 정의 그 자체(leftRecord.key == rightRecord.key)이므로
요건 목록에 넣지 않습니다.
| # | 요건 | Streams가 검증하는가 | 위반 시 무슨 일이 일어나는가 |
|---|---|---|---|
| ① | 파티션 수가 같아야 합니다 (조인 양쪽 입력 토픽) | 검증합니다 파티션 할당 단계, 즉 런타임에 |
TopologyException(런타임 예외)이 발생해 즉시 실패합니다. 알아차리기 쉽습니다 |
| ② | 입력 토픽에 쓰는 모든 애플리케이션이 같은 파티셔너를 써야 합니다. Producer API는 partitioner.class, Streams는 StreamPartitioner(예: KStream#to()) |
검증하지 못합니다 공식 문서가 "사용자 책임"이라고 명시 |
예외 없이 조용히 잘못된 결과가 나옵니다. 같은 키가 다른 파티션에 들어가 조인이 매칭되지 않습니다 — 가장 찾기 어려운 버그 |
# 1) 양쪽 파티션 수 확인
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic customers
// KStream 이 적은 쪽이면
KStream<String, Order> repartitioned =
orders.repartition(Repartitioned.numberOfPartitions(24));
// KTable 이 적은 쪽이면 — toStream → repartition → toTable
KTable<String, Customer> repartitionedTable =
customers.toStream()
.repartition(Repartitioned.numberOfPartitions(24))
.toTable();
// stream-table 조인이라면 KStream 쪽을 리파티션하는 것이 권장됩니다.
// KTable 을 리파티션하면 상태 저장소가 하나 더 생길 수 있습니다.
내부 토픽
| 토픽 종류 | cleanup.policy |
보관 | 용도 |
|---|---|---|---|
| repartition 토픽 | delete |
retention -1 (무한) |
키를 다시 분배하기 위한 중간 토픽. Streams가 처리 후 직접 purge합니다 (repartition.purge.interval.ms) |
| changelog (key-value 스토어) | compact |
컴팩션으로 최신 값만 | 상태 저장소 복구용. 이 토픽을 지우면 상태를 복구할 수 없습니다 |
| changelog (윈도 스토어) | delete,compact |
24시간 + 윈도 스토어 설정값 | 윈도는 만료되므로 삭제와 컴팩션을 함께 씁니다 |
| changelog (versioned 스토어) | compact |
min.compaction.lag.ms = 24시간 + historyRetentionMs |
과거 버전 조회를 보장하기 위해 컴팩션을 지연시킵니다 |
| 모든 내부 토픽 | — | — | message.timestamp.type=CreateTime으로 설정됩니다 |
핵심 설정
| 설정 | 기본값 | 역할과 튜닝 포인트 |
|---|---|---|
application.id | 필수 | 컨슈머 group.id와 내부 토픽 접두어로 쓰입니다. 바꾸면 완전히 새 애플리케이션이 되어 상태와 오프셋이 초기화됩니다 |
bootstrap.servers | 필수 | 2개 이상 지정하세요 |
num.stream.threads | 1 | 인스턴스당 스레드 수. 전체 병렬성 상한은 입력 토픽의 파티션 수입니다. 그 이상 늘리면 유휴 스레드만 생깁니다 |
processing.guarantee | at_least_once | exactly_once_v2로 바꿔야 EOS가 켜집니다. 기본값이 아닙니다. 켜면 commit.interval.ms가 종단 지연의 하한이 됩니다 |
replication.factor | -1 | -1은 브로커 기본값(default.replication.factor, 기본 1)을 씁니다. 운영에서는 3으로 명시하세요 |
state.dir | ${java.io.tmpdir} | RocksDB 상태 저장소 위치. 컨테이너에서 임시 볼륨을 쓰면 재시작마다 changelog를 전부 다시 읽습니다 |
num.standby.replicas | 0 | 1로 두면 hot standby가 상태를 미리 따라가 페일오버 복구 시간이 극적으로 줄어듭니다. 대신 브로커 트래픽과 디스크가 늘어납니다 |
max.warmup.replicas | 2 | 동시에 워밍업할 태스크 수. 스케일 아웃 속도와 복제 부하의 트레이드오프 |
acceptable.recovery.lag | 10000 | 이 lag 이내면 "따라잡았다"고 보고 활성 태스크로 승격합니다 |
probing.rebalance.interval.ms | 600000 (10분) | 워밍업 진행을 확인하는 리밸런스 주기 |
commit.interval.ms | 30000 (30초) | 오프셋·상태 커밋 주기. EOS에서는 이 값이 종단 지연을 지배합니다. 낮추면 지연이 줄고 트랜잭션 수가 늘어납니다 |
statestore.cache.max.bytes | 10485760 (10MB) | 인스턴스 전체 상태 저장소 캐시. 크게 하면 중간 결과 방출이 줄어 downstream 부하가 감소하고, 결과가 늦게 보입니다 |
cache.max.bytes.buffering | 10485760 | 같은 값의 이전 이름입니다. 신규 코드에서는 statestore.cache.max.bytes를 쓰세요 |
topology.optimization | none | all로 두면 불필요한 repartition 토픽을 줄입니다. 기존 앱에 켜면 토폴로지가 바뀌어 호환 문제가 생길 수 있습니다 |
max.task.idle.ms | 0 | 한쪽 입력이 비었을 때 조인·병합에서 기다릴 시간. 올리면 이벤트 시간 정렬이 좋아지고 지연이 늘어납니다 |
task.timeout.ms | 300000 (5분) | 재시도 가능한 오류로 태스크가 멈춰 있을 수 있는 시간 |
poll.ms | 100 | 내부 컨슈머 poll() 블록 시간 |
buffered.records.per.partition | 1000 | 파티션당 버퍼링 레코드 수. 메모리와 이벤트 시간 정렬에 영향 |
default.timestamp.extractor | org.apache.kafka.streams.processor.FailOnInvalidTimestamp | 기본값은 잘못된 타임스탬프에서 실패합니다. 레거시 데이터를 다루면 교체가 필요할 수 있습니다 |
deserialization.exception.handler | org.apache.kafka.streams.errors.LogAndFailExceptionHandler | 기본은 로그 후 실패입니다. 포이즌 메시지 하나로 앱이 멈추므로 DLQ 전략과 함께 설계하세요 |
production.exception.handler | org.apache.kafka.streams.errors.DefaultProductionExceptionHandler | 출력 실패 처리 |
processing.exception.handler | org.apache.kafka.streams.errors.LogAndFailProcessingExceptionHandler | 처리 로직 예외. 기본은 로그 후 실패 |
errors.dead.letter.queue.topic.name | null | 설정하지 않으면 DLQ가 동작하지 않습니다 |
default.dsl.store | rocksDB | in_memory로 바꾸면 디스크를 안 쓰지만 힙이 커집니다 |
rocksdb.config.setter | null | 블록 캐시·write buffer 튜닝 지점. 메모리 폭주의 흔한 원인이기도 합니다 |
rack.aware.assignment.strategy | none | 멀티 AZ에서 크로스 AZ 트래픽 비용을 줄입니다. rack.aware.assignment.tags와 함께 씁니다 |
application.server | "" | 인터랙티브 쿼리용 host:port. 설정하면 다른 인스턴스로 쿼리를 라우팅할 수 있습니다 |
group.protocol | classic | Streams Rebalance Protocol을 쓰려면 streams로 설정합니다. 클라이언트 기본값은 classic입니다 — 아래 버전 노트 참조 |
upgrade.from | null | 메이저 업그레이드 시 2단계 롤링에 필요합니다 |
windowstore.changelog.additional.retention.ms | 86400000 (1일) | 윈도 스토어 changelog에 더할 추가 보관 시간. 내부 토픽 표의 "24시간"이 이 값입니다 |
repartition.purge.interval.ms | 30000 (30초) | repartition 토픽의 처리 완료 구간을 삭제 요청하는 주기 |
state.cleanup.delay.ms | 600000 (10분) | 태스크가 다른 인스턴스로 옮겨간 뒤 로컬 상태를 지우기까지 대기하는 시간 |
request.timeout.ms | 40000 (40초) | Streams 내부 클라이언트의 요청 타임아웃. 일반 클라이언트 기본값(30초)과 다릅니다 |
Streams Rebalance Protocol
| 설정 | 기본값 | 역할 |
|---|---|---|
group.streams.session.timeout.ms | 45000 | Streams 그룹 멤버의 세션 타임아웃 |
group.streams.heartbeat.interval.ms | 5000 | 멤버에게 부여되는 하트비트 간격 |
group.streams.num.standby.replicas | 0 | 태스크별 standby 수의 브로커 측 기본값 |
group.streams.max.standby.replicas | 2 | 허용 상한 |
group.streams.initial.rebalance.delay.ms | 3000 | 첫 리밸런스 지연 |
group.streams.assignment.interval.ms | 1000 | 브로커가 할당을 재계산하는 주기 |
group.coordinator.rebalance.protocols | classic,consumer,streams | streams가 이 목록에 있어야 프로토콜이 동작합니다 |
운영 명령
export BS=localhost:9092
# Streams 그룹 목록
bin/kafka-streams-groups.sh --bootstrap-server $BS --list
# 내부 토픽 포함 전체 확인
bin/kafka-topics.sh --bootstrap-server $BS --list | grep "^my-streams-app-"
# lag 확인 — Streams 도 결국 컨슈머 그룹입니다 (application.id = group.id)
bin/kafka-consumer-groups.sh --bootstrap-server $BS --describe --group my-streams-app
# 앱 완전 초기화:
# 1) 애플리케이션 인스턴스를 모두 내립니다
# 2) 아래 도구로 오프셋과 내부 토픽을 정리합니다
bin/kafka-streams-application-reset.sh --bootstrap-server $BS \
--application-id my-streams-app \
--input-topics orders
# 3) 로컬 상태 디렉터리도 직접 지워야 합니다 (도구가 지우지 않습니다)
rm -rf /var/lib/kafka-streams/my-streams-app
자주 걸리는 지점
시험 포인트 정리
이어서 볼 곳
공식 문서 출처
타입 전이·리파티션 유발 여부·윈도 정의·co-partitioning 요구사항·내부 토픽 설정은
Apache Kafka 4.3.1 공식 문서에서 확인했습니다.
설정 기본값은 공식 Streams 설정 표에서 확인했습니다.
toTable()의 리파티션 동작은 문서에 명시가 없어 표에 그대로 적었습니다.
- Streams DSL — 연산자 표, 리파티션 마킹 규칙, 윈도 4종, 조인 5종, co-partitioning 요구사항
- Streams Configs — 설정 기본값
- Kafka Streams Configs (설정 표)
- Managing Streams Application Topics — 내부 토픽 명명 규칙과 기본 설정
- Streams Core Concepts — 태스크·병렬성·시간 개념
- Streams Architecture — 스레드·태스크 모델
- Streams Rebalance Protocol — 4.2부터 새 클러스터 기본 활성,
group.protocol=streams,group.streams.*브로커 설정 - Application Reset Tool — 리셋 절차와 로컬 상태 삭제