3초 요약

DSL 연산자 표

Kafka Streams DSL 연산자 — 타입 전이는 공식 문서의 Transformation 표 기준
연산 입력 → 출력 상태 저장 리파티션 유발 메모
소스 · 싱크
stream() 입력 토픽 → KStream 무상태 안 함 레코드 스트림으로 읽습니다 (append-only 의미론)
table() 입력 토픽 → KTable 상태 안 함 changelog 의미론. 같은 키의 새 값이 이전 값을 덮어씁니다
globalTable() 입력 토픽 → GlobalKTable 상태 안 함 모든 인스턴스가 전체 파티션의 사본을 갖습니다. 그래서 co-partitioning이 불필요합니다
to() KStream → void 조건부 출력 토픽 파티션 수가 다르거나, 스트림이 마킹됐거나, 커스텀 StreamPartitioner를 주거나, 키가 null이면 리파티션합니다
무상태 변환 (stateless)
filter() KStreamKStream, KTableKTable 무상태 안 함 조건이 참인 레코드만 남깁니다
filterNot() KStreamKStream, KTableKTable 무상태 안 함 조건이 거짓인 레코드만 남깁니다
map() KStreamKStream 무상태 마킹함 키를 바꿀 수 있으므로 리파티션 대상이 됩니다. 값만 바꾸면 mapValues를 쓰세요
mapValues() KStreamKStream, KTableKTable 무상태 안 함 키를 보존하므로 리파티션이 필요 없습니다. map보다 항상 우선
flatMap() KStreamKStream 무상태 마킹함 0..N개 레코드로 확장. 키 변경 가능
flatMapValues() KStreamKStream 무상태 안 함 키를 보존한 채 값을 확장. flatMap보다 우선
selectKey() KStreamKStream 무상태 마킹함 키를 명시적으로 바꿉니다. 조인 전 co-partitioning을 맞출 때 씁니다
split() / Branch KStreamBranchedKStream 무상태 안 함 술어를 순서대로 평가해 첫 매치 하나의 스트림으로 보냅니다. 매치가 없으면 default 브랜치 또는 폐기
merge() KStreamKStream 무상태 안 함 두 스트림을 합칩니다. 순서는 보장되지 않습니다
peek() KStreamKStream 무상태 안 함 부수효과를 만들면서 스트림을 그대로 통과시킵니다 (로깅·메트릭용)
foreach() KStream→void, KTable→void 무상태 안 함 종료 연산입니다. 이후 처리를 이어갈 수 없습니다 — 이어가려면 peek
print() KStream→void 무상태 안 함 디버깅용 종료 연산
toStream() KTableKStream 무상태 안 함 테이블을 changelog 스트림으로 봅니다
toTable() KStreamKTable 상태 문서에 명시 없음 스트림을 테이블로 변환합니다. 앞선 연산이 키를 바꿨다면 그 마킹이 여기까지 이어집니다
repartition() KStreamKStream 무상태 항상 유발 파티션 수를 명시적으로 지정할 수 있습니다. 자동 리파티션이 걸리지 않는 process() 계열 앞에서 필요합니다
그룹화 (집계 전 단계)
groupByKey() KStreamKGroupedStream 무상태 마킹된 경우에만 groupBy보다 항상 우선. 이미 원하는 키라면 리파티션이 전혀 일어나지 않습니다
groupBy() KStreamKGroupedStream, KTableKGroupedTable 무상태 항상 유발 키를 새로 계산하므로 무조건 리파티션합니다. 가능하면 groupByKey로 대체하세요
cogroup() KGroupedStreamCogroupedKStream 안 함 이미 그룹화된 스트림을 받으므로 리파티션이 그 단계에서 끝나 있습니다
집계 (상태 저장)
count() KGroupedStreamKTable, KGroupedTableKTable 상태 그룹화 단계에서 결정 키별 레코드 수
reduce() KGroupedStreamKTable, KGroupedTableKTable 상태 그룹화 단계에서 결정 결과 타입이 입력 값 타입과 같아야 합니다
aggregate() KGroupedStreamKTable, KGroupedTableKTable 상태 그룹화 단계에서 결정 결과 타입을 바꿀 수 있습니다. initializer가 필요합니다
windowedBy() KGroupedStreamTimeWindowedKStream / SessionWindowedKStream 상태 그룹화 단계에서 결정 윈도 스토어를 씁니다. 기본 보관 기간은 1일이며 Materialized#withRetention()으로 변경합니다
suppress() KTableKTable 상태 안 함 윈도의 최종 결과만 내보냅니다. 중간 갱신을 억제해 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() KStreamKStream 스토어 연결 시 상태 자동 리파티션 안 함 앞에서 키를 바꿨다면 repartition()을 직접 호출해야 합니다. 이것이 process()의 가장 큰 함정
processValues() KStreamKStream 스토어 연결 시 상태 안 함 키를 바꿀 수 없는(FixedKeyProcessor) 버전이라 리파티션이 필요 없습니다

map이 키를 바꿀 수 있다고 판단해 스트림이 마킹되고, groupByKey에서 리파티션 토픽이 생깁니다. 값만 바꾸는데도 네트워크 왕복이 추가됩니다.

Topology.java
builder.stream("orders", Consumed.with(Serdes.String(), orderSerde))
       // 키를 그대로 돌려주더라도 Streams 는 알 수 없습니다.
       // map 이 호출된 것만으로 "키가 바뀔 수 있다"고 마킹됩니다.
       .map((k, v) -> KeyValue.pair(k, v.withTax()))
       .groupByKey()          // ← 여기서 리파티션 토픽 생성
       .count();

mapValues는 키를 바꿀 수 없으므로 마킹되지 않고, groupByKey가 리파티션 없이 진행됩니다.

Topology.java
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종

윈도 타입 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 안의 이벤트는 기존 세션에 병합되며, 늦게 온 레코드가 두 세션을 하나로 합칠 수도 있습니다
4종 정의 코드 — 공식 문서 예시 기준
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));

조인 매트릭스

조인 5종 — 윈도 필요 여부, co-partitioning 요구, 지원 조인 타입
조인 윈도 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)이므로 요건 목록에 넣지 않습니다.

co-partitioning 요건 2개 — Streams가 검증하는 것과 하지 않는 것을 구분하는 것이 실무의 핵심입니다
# 요건 Streams가 검증하는가 위반 시 무슨 일이 일어나는가
파티션 수가 같아야 합니다 (조인 양쪽 입력 토픽) 검증합니다
파티션 할당 단계, 즉 런타임
TopologyException(런타임 예외)이 발생해 즉시 실패합니다. 알아차리기 쉽습니다
입력 토픽에 쓰는 모든 애플리케이션이 같은 파티셔너를 써야 합니다. Producer API는 partitioner.class, Streams는 StreamPartitioner(예: KStream#to()) 검증하지 못합니다
공식 문서가 "사용자 책임"이라고 명시
예외 없이 조용히 잘못된 결과가 나옵니다. 같은 키가 다른 파티션에 들어가 조인이 매칭되지 않습니다 — 가장 찾기 어려운 버그
co-partitioning을 맞추는 절차 — 공식 문서 권고 순서
# 1) 양쪽 파티션 수 확인
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders
bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic customers
2) 파티션 수가 적은 쪽을 리파티션 — 병목을 피하려면 큰 쪽에 맞춥니다
// 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으로 설정됩니다

핵심 설정

Kafka Streams 핵심 설정 33개 — Apache Kafka 4.3.1 공식 문서 기준 기본값
설정 기본값 역할과 튜닝 포인트
application.id필수컨슈머 group.id와 내부 토픽 접두어로 쓰입니다. 바꾸면 완전히 새 애플리케이션이 되어 상태와 오프셋이 초기화됩니다
bootstrap.servers필수2개 이상 지정하세요
num.stream.threads1인스턴스당 스레드 수. 전체 병렬성 상한은 입력 토픽의 파티션 수입니다. 그 이상 늘리면 유휴 스레드만 생깁니다
processing.guaranteeat_least_onceexactly_once_v2로 바꿔야 EOS가 켜집니다. 기본값이 아닙니다. 켜면 commit.interval.ms가 종단 지연의 하한이 됩니다
replication.factor-1-1은 브로커 기본값(default.replication.factor, 기본 1)을 씁니다. 운영에서는 3으로 명시하세요
state.dir${java.io.tmpdir}RocksDB 상태 저장소 위치. 컨테이너에서 임시 볼륨을 쓰면 재시작마다 changelog를 전부 다시 읽습니다
num.standby.replicas01로 두면 hot standby가 상태를 미리 따라가 페일오버 복구 시간이 극적으로 줄어듭니다. 대신 브로커 트래픽과 디스크가 늘어납니다
max.warmup.replicas2동시에 워밍업할 태스크 수. 스케일 아웃 속도와 복제 부하의 트레이드오프
acceptable.recovery.lag10000이 lag 이내면 "따라잡았다"고 보고 활성 태스크로 승격합니다
probing.rebalance.interval.ms600000 (10분)워밍업 진행을 확인하는 리밸런스 주기
commit.interval.ms30000 (30초)오프셋·상태 커밋 주기. EOS에서는 이 값이 종단 지연을 지배합니다. 낮추면 지연이 줄고 트랜잭션 수가 늘어납니다
statestore.cache.max.bytes10485760 (10MB)인스턴스 전체 상태 저장소 캐시. 크게 하면 중간 결과 방출이 줄어 downstream 부하가 감소하고, 결과가 늦게 보입니다
cache.max.bytes.buffering10485760같은 값의 이전 이름입니다. 신규 코드에서는 statestore.cache.max.bytes를 쓰세요
topology.optimizationnoneall로 두면 불필요한 repartition 토픽을 줄입니다. 기존 앱에 켜면 토폴로지가 바뀌어 호환 문제가 생길 수 있습니다
max.task.idle.ms0한쪽 입력이 비었을 때 조인·병합에서 기다릴 시간. 올리면 이벤트 시간 정렬이 좋아지고 지연이 늘어납니다
task.timeout.ms300000 (5분)재시도 가능한 오류로 태스크가 멈춰 있을 수 있는 시간
poll.ms100내부 컨슈머 poll() 블록 시간
buffered.records.per.partition1000파티션당 버퍼링 레코드 수. 메모리와 이벤트 시간 정렬에 영향
default.timestamp.extractororg.apache.kafka.streams.processor.FailOnInvalidTimestamp기본값은 잘못된 타임스탬프에서 실패합니다. 레거시 데이터를 다루면 교체가 필요할 수 있습니다
deserialization.exception.handlerorg.apache.kafka.streams.errors.LogAndFailExceptionHandler기본은 로그 후 실패입니다. 포이즌 메시지 하나로 앱이 멈추므로 DLQ 전략과 함께 설계하세요
production.exception.handlerorg.apache.kafka.streams.errors.DefaultProductionExceptionHandler출력 실패 처리
processing.exception.handlerorg.apache.kafka.streams.errors.LogAndFailProcessingExceptionHandler처리 로직 예외. 기본은 로그 후 실패
errors.dead.letter.queue.topic.namenull설정하지 않으면 DLQ가 동작하지 않습니다
default.dsl.storerocksDBin_memory로 바꾸면 디스크를 안 쓰지만 힙이 커집니다
rocksdb.config.setternull블록 캐시·write buffer 튜닝 지점. 메모리 폭주의 흔한 원인이기도 합니다
rack.aware.assignment.strategynone멀티 AZ에서 크로스 AZ 트래픽 비용을 줄입니다. rack.aware.assignment.tags와 함께 씁니다
application.server""인터랙티브 쿼리용 host:port. 설정하면 다른 인스턴스로 쿼리를 라우팅할 수 있습니다
group.protocolclassicStreams Rebalance Protocol을 쓰려면 streams로 설정합니다. 클라이언트 기본값은 classic입니다 — 아래 버전 노트 참조
upgrade.fromnull메이저 업그레이드 시 2단계 롤링에 필요합니다
windowstore.changelog.additional.retention.ms86400000 (1일)윈도 스토어 changelog에 더할 추가 보관 시간. 내부 토픽 표의 "24시간"이 이 값입니다
repartition.purge.interval.ms30000 (30초)repartition 토픽의 처리 완료 구간을 삭제 요청하는 주기
state.cleanup.delay.ms600000 (10분)태스크가 다른 인스턴스로 옮겨간 뒤 로컬 상태를 지우기까지 대기하는 시간
request.timeout.ms40000 (40초)Streams 내부 클라이언트의 요청 타임아웃. 일반 클라이언트 기본값(30초)과 다릅니다

Streams Rebalance Protocol

Streams 그룹 관련 브로커 설정 — 클라이언트 값이 아니라 브로커가 관리합니다
설정기본값역할
group.streams.session.timeout.ms45000Streams 그룹 멤버의 세션 타임아웃
group.streams.heartbeat.interval.ms5000멤버에게 부여되는 하트비트 간격
group.streams.num.standby.replicas0태스크별 standby 수의 브로커 측 기본값
group.streams.max.standby.replicas2허용 상한
group.streams.initial.rebalance.delay.ms3000첫 리밸런스 지연
group.streams.assignment.interval.ms1000브로커가 할당을 재계산하는 주기
group.coordinator.rebalance.protocolsclassic,consumer,streamsstreams가 이 목록에 있어야 프로토콜이 동작합니다

운영 명령

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()의 리파티션 동작은 문서에 명시가 없어 표에 그대로 적었습니다.