기본개념 · 10장
Kafka Streams와 ksqlDB
Kafka Streams는 별도 클러스터를 세우지 않고 애플리케이션 안에서 돌아가는 라이브러리로
스트림 처리를 합니다. 이 장에서는 토폴로지와 태스크의 관계, KStream/KTable/GlobalKTable의 의미론 차이,
그리고 시험과 실무 양쪽에서 가장 자주 물리는 두 가지 —
map은 리파티션을 유발하고 mapValues는 유발하지 않는다는 사실과
조인의 co-partitioning 요구사항 — 을 정확히 정리합니다.
Kafka Streams는 CCDAK에서 12%, 함께 다루는 테스팅은 8% 비중입니다.
학습 목표
- Streams가 왜 클러스터가 아니라 라이브러리인지, 그것이 배포와 확장에 무엇을 뜻하는지 설명할 수 있습니다.
- 같은 입력에 대해 KStream · KTable · GlobalKTable이 서로 다른 결과를 내는 이유를 설명할 수 있습니다.
- 어떤 DSL 연산이 상태를 만들고 어떤 연산이 리파티션을 유발하는지 구분할 수 있습니다.
- 조인 4종의 윈도우 필요 여부와 co-partitioning 요구사항을 판단할 수 있습니다.
- tumbling · hopping · sliding · session 윈도우를 구분하고 grace period의 역할을 설명할 수 있습니다.
TopologyTestDriver로 브로커 없이 토폴로지를 테스트할 수 있습니다.
클러스터가 아니라 라이브러리입니다
스트림 처리 프레임워크는 보통 자체 클러스터와 리소스 매니저를 요구합니다. Kafka Streams는 그 방향을 택하지 않았습니다. 공식 문서의 표현대로 "어떤 Java 애플리케이션에도 쉽게 임베드할 수 있는 단순하고 가벼운 클라이언트 라이브러리"이고, 내부 메시징 계층으로 Apache Kafka 자체 외에 어떤 외부 시스템에도 의존하지 않습니다. 문서는 "Kafka Streams는 리소스 매니저가 아니며, 스트림 처리 애플리케이션이 도는 곳이면 어디서든 돈다"고 명시합니다.
| 관점 | Consumer + Producer 직접 구현 | Kafka Streams |
|---|---|---|
| 배포 단위 | 일반 애플리케이션 | 일반 애플리케이션 (동일). 별도 클러스터가 없습니다 |
| 상태 관리 | 직접 구현 (외부 DB 등) | 로컬 상태 저장소 + changelog 토픽을 프레임워크가 관리 |
| 장애 복구 | 직접 구현 | 태스크가 다른 인스턴스에서 자동 재시작되고 changelog로 상태를 복원 |
| exactly-once | 트랜잭션 API를 직접 조합 | processing.guarantee=exactly_once_v2 한 줄 |
| 시간 처리 | 직접 구현 | event time 기반 윈도우, 순서 어긋난 데이터 처리 내장 |
| 확장 방법 | 컨슈머 인스턴스 추가 | 인스턴스 또는 num.stream.threads 추가. 파티션 재분배는 자동 |
| 적합한 일 | 단순 소비·발행, 외부 시스템 연동 | 변환·집계·조인·윈도우처럼 상태와 시간이 필요한 처리 |
토폴로지 — source · processor · sink
Streams 애플리케이션은 프로세서 토폴로지(processor topology)로 계산 로직을 정의합니다. 토폴로지는 스트림 프로세서(노드)를 스트림(엣지)으로 연결한 그래프이며, 두 개의 특별한 프로세서가 있습니다.
- source processor — 상위 프로세서가 없는 노드. 하나 이상의 Kafka 토픽에서 레코드를 읽어 하위로 전달합니다.
- sink processor — 하위 프로세서가 없는 노드. 받은 레코드를 지정한 Kafka 토픽으로 보냅니다.
토폴로지를 정의하는 방법은 두 가지입니다 —
Kafka Streams DSL(map, filter, join, 집계 등을 바로 제공)과
더 저수준인 Processor API(커스텀 프로세서를 직접 정의하고 상태 저장소에 접근).
토폴로지는 논리적 추상일 뿐이고, 런타임에는 병렬 처리를 위해 인스턴스화되고 복제됩니다.
import java.util.Arrays;
import java.util.Locale;
import java.util.Properties;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
public final class WordCount {
public static Topology buildTopology() {
StreamsBuilder builder = new StreamsBuilder();
// source processor: 입력 토픽
KStream<String, String> lines = builder.stream("text-lines");
KTable<String, Long> counts = lines
// 값만 바꿉니다 → 리파티션을 유발하지 않습니다
.flatMapValues(line -> Arrays.asList(line.toLowerCase(Locale.ROOT).split("\\W+")))
// 키를 단어로 바꿉니다 → 여기서 리파티션이 필요해집니다
.groupBy((key, word) -> word)
// stateful: 상태 저장소와 changelog 토픽이 생깁니다
.count(Materialized.as("counts-store"));
// sink processor: 출력 토픽
counts.toStream().to("word-counts", Produced.with(Serdes.String(), Serdes.Long()));
return builder.build();
}
public static void main(String[] args) {
Properties props = new Properties();
// application.id 는 컨슈머 group.id, client.id 접두어, changelog 토픽 접두어로 함께 쓰입니다
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// changelog · 리파티션 토픽의 복제 계수. 기본값 -1 은 브로커 기본값을 따릅니다
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);
KafkaStreams streams = new KafkaStreams(buildTopology(), props);
// 예외로 스트림 스레드가 죽었을 때의 정책을 반드시 정하세요.
// 정하지 않으면 스레드가 하나씩 사라져 처리 용량이 조용히 줄어듭니다.
streams.setUncaughtExceptionHandler(
exception -> StreamThreadExceptionResponse.REPLACE_THREAD);
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
streams.start();
}
}
KStream · KTable · GlobalKTable
같은 토픽을 읽어도 무엇으로 해석하느냐에 따라 결과가 달라집니다.
공식 문서가 드는 예시가 가장 명확합니다 — ("alice", 1) 다음에 ("alice", 3)이 들어올 때
사용자별 값을 합산하면,
- KStream이면 4입니다. 두 번째 레코드는 앞 레코드를 대체하지 않는 별개의 사실이기 때문입니다.
- KTable이면 3입니다. 두 번째 레코드가 앞 값의 갱신으로 해석되기 때문입니다.
| KStream | KTable | GlobalKTable | |
|---|---|---|---|
| 해석 | record stream — 각 레코드가 독립적인 사실 | changelog stream — 각 레코드가 갱신 | changelog stream (KTable과 동일) |
| 같은 키 반복 | INSERT — append-only 원장에 추가 | UPSERT — 기존 행을 덮어씀 | UPSERT |
null value |
그냥 값이 null인 레코드 | DELETE (tombstone) | DELETE (tombstone) |
| 인스턴스가 채우는 데이터 | 배정받은 파티션 | 배정받은 파티션만 (5파티션·5인스턴스면 각 1개) | 모든 파티션 (각 인스턴스가 전체 사본을 가짐) |
| 시간 개념 | 있음 | 있음 | 없습니다 — 공식 문서가 명시적으로 대비합니다 |
| 조인 시 co-partitioning | 필요 | 필요 (equi-join) | 불필요 |
| 키가 아닌 값으로 조인 | 불가 | foreign-key join으로 가능 | 가능 (KeyValueMapper) |
| 비용 | — | — | 로컬 저장 공간과 브로커 부하가 큼 (토픽 전체를 각 인스턴스가 읽음) |
| 어떤 데이터에 쓰나 | 클릭, 결제 트랜잭션, 센서 측정, 로그 | 사용자 프로필, 재고 수량, 계정 잔액 | 국가 코드, 환율처럼 작고 자주 변하지 않는 참조 데이터 |
("alice",1) → ("alice",3)을
KStream(합 4) · KTable(합 3) · GlobalKTable(전 파티션 복제)로 읽었을 때의 결과 차이
stateless와 stateful, 그리고 리파티션
DSL 연산을 분류하는 축은 두 개이고 서로 독립입니다 —
상태 저장소를 만드는가와 리파티션을 유발하는가입니다.
map은 stateless이지만 리파티션을 유발하고, count는 stateful이지만
그 자체가 리파티션을 유발하지는 않습니다. 이 두 축을 섞으면 문제를 틀립니다.
| 연산 | 상태 | 리파티션 | 비고 |
|---|---|---|---|
filter / filterNot | stateless | 유발 안 함 | 키를 바꾸지 않습니다 |
mapValues | stateless | 유발 안 함 | 값만 바꿉니다 |
flatMapValues | stateless | 유발 안 함 | 값만 바꾸며 개수를 늘립니다 |
map | stateless | 표시함 | 키를 바꿀 수 있으므로 리파티션 대상으로 표시(mark)됩니다 |
flatMap | stateless | 표시함 | 가능하면 flatMapValues를 쓰라고 문서가 권합니다 |
selectKey | stateless | 표시함 | 키를 바꾸는 것이 목적입니다 |
groupByKey | stateless (그룹화만) | 표시됐을 때만 | 스트림이 리파티션 대상으로 표시된 경우에만 실제로 리파티션합니다 |
groupBy | stateless (그룹화만) | 항상 | 키를 새로 정하므로 언제나 리파티션합니다. 가능하면 groupByKey를 쓰세요 |
repartition | stateless | 항상 | 파티션 수를 지정해 명시적으로 재분배. 내부 토픽은 Streams가 관리·자동 삭제합니다 |
branch / merge / peek / foreach / print | stateless | 유발 안 함 | foreach·print는 종단 연산입니다 |
count / reduce / aggregate | stateful | 앞 연산에 달림 | 상태 저장소 + changelog 토픽이 생깁니다 |
| 모든 윈도우 연산 | stateful | 앞 연산에 달림 | 윈도우 상태 저장소를 씁니다 |
| 모든 조인 | stateful | 표시됐을 때만 | 양쪽이 표시됐으면 양쪽이 리파티션됩니다 |
suppress | stateful | 유발 안 함 | 윈도우 최종 결과만 내보낼 때 씁니다 |
cogroup | stateful | 유발 안 함 | 입력이 이미 그룹화되어 있는 것이 전제이므로 리파티션을 유발하지 않습니다 |
map은 유발하고 mapValues는 유발하지 않는다는 대비가 중심
상태 저장소와 changelog
stateful 연산을 쓰면 DSL이 상태 저장소(state store)를 자동으로 만들고 관리합니다. 기본 구현은 디스크 기반 키-값 저장소(RocksDB)이고, 인메모리 해시맵도 선택할 수 있습니다. 내고장성은 changelog 토픽이 담당합니다.
- 각 상태 저장소마다 복제되는 changelog Kafka 토픽이 유지되며 모든 상태 갱신이 기록됩니다.
- changelog도 파티션되어 있어 각 저장소 인스턴스(=태스크)마다 전용 changelog 파티션을 가집니다.
- changelog 토픽에는 로그 컴팩션이 켜져 있어 옛 데이터가 안전하게 정리됩니다 — 무한히 커지지 않습니다.
- 태스크가 다른 머신에서 재시작되면 changelog를 재생(replay)해 장애 이전 상태로 복원한 뒤 처리를 이어갑니다.
문제는 복원 시간입니다. 태스크 재초기화 비용은 대부분 changelog 재생 시간이 차지합니다.
이를 줄이는 장치가 standby replica입니다.
num.standby.replicas를 올리면 상태의 완전 복제본을 다른 인스턴스에 미리 유지하고,
태스크가 이동할 때 이미 최신 상태를 가진 인스턴스로 배정합니다.
2.6부터는 완전히 따라잡은 로컬 복사본을 가진 인스턴스가 존재한다면 반드시 그 인스턴스에 배정됩니다.
시간 — event time · processing time · ingestion time
| 개념 | 언제의 시각인가 | 누가 부여하는가 |
|---|---|---|
| event time | 이벤트가 소스에서 실제로 발생한 시점 | 이벤트를 만든 쪽 (예: GPS 센서가 위치 변화를 감지한 순간) |
| processing time | 스트림 처리 애플리케이션이 그 레코드를 처리한 시점 | 처리 애플리케이션. event time보다 밀리초~수 시간 늦을 수 있습니다 |
| ingestion time | 브로커가 토픽 파티션에 저장한 시점 | 브로커. 레코드가 끝까지 처리되지 않아도 ingestion time은 존재합니다 |
Streams는 TimestampExtractor 인터페이스로 모든 레코드에 타임스탬프를 부여합니다.
이 값이 진행 상황을 나타내는 stream time이고, 실제 실행 시각인 wall-clock time과 구분됩니다.
stream time은 새 레코드가 프로세서에 도착할 때만 전진합니다 —
데이터가 멈추면 시간도 멈춥니다. 윈도우가 닫히지 않는 현상의 원인이 대부분 이것입니다.
| 설정 | 기본값 | 설명 | 튜닝 포인트 |
|---|---|---|---|
default.timestamp.extractor |
FailOnInvalidTimestamp |
embed된 타임스탬프를 그대로 쓰고, 유효하지 않으면 예외를 던집니다 | 구버전 프로듀서가 만든 음수 타임스탬프가 섞여 있으면 애플리케이션이 죽습니다. 처리 시각으로 대체하려면 WallclockTimestampExtractor를 쓰지만, event time 의미론을 포기하는 것임을 알아야 합니다. |
max.task.idle.ms |
0 | 일부 입력 파티션만 데이터가 있을 때, 나머지를 기다릴 최대 시간 | 기본값 0은 기다리지 않습니다 → 조인·merge에서 순서가 어긋난 결과가 나올 수 있습니다. 조인 정확도가 중요하면 올립니다. |
default.deserialization.exception.handler |
LogAndFailExceptionHandler |
역직렬화 실패 시 로그를 남기고 실패 | 독성 레코드(poison pill) 하나로 애플리케이션이 멈춥니다. 건너뛰려면 다른 핸들러가 필요합니다. |
processing.exception.handler |
LogAndFailProcessingExceptionHandler |
처리 중 예외 시 로그를 남기고 실패. 글로벌 상태 저장소 갱신에는 적용되지 않습니다 | 글로벌 스레드의 예외는 uncaught exception handler로 올라갑니다. |
출력 레코드의 타임스탬프는 어떻게 정해지는가
| 상황 | 출력 타임스탬프 |
|---|---|
입력 레코드를 처리해 출력 (context.forward() 등) | 입력 레코드의 타임스탬프를 그대로 상속 |
주기 함수(Punctuator#punctuate())에서 출력 | 스트림 태스크의 현재 내부 시각 |
| stream-stream / table-table 조인 | max(left.ts, right.ts) |
| stream-table 조인 | stream 쪽 레코드의 타임스탬프 |
| 집계 | 키별(또는 윈도우별)로 기여한 모든 입력 중 최대 타임스탬프 |
| stateless 연산 | 입력 타임스탬프 통과. flatMap류는 모든 출력이 같은 값을 상속 |
윈도우 4종
윈도우는 stateful 연산을 위해 같은 키를 가진 레코드를 시간 구간으로 다시 묶는 장치입니다. 공식 문서가 강조하는 전제: 윈도우는 레코드 키별로 따로 추적됩니다.
| 종류 | 정의 요소 | 겹침 | 정렬 기준 | API |
|---|---|---|---|---|
| tumbling | 크기(size)만 | 겹치지 않음. 레코드는 정확히 하나의 윈도우에 속합니다 | epoch 정렬 — 5000ms면 [0,5000) [5000,10000). 하한 포함·상한 제외 |
TimeWindows.ofSizeAndGrace(size, grace) |
| hopping | 크기 + advance(hop) | 겹칩니다. 한 레코드가 여러 윈도우에 속할 수 있습니다 | epoch 정렬 — size 5000·hop 3000이면 [0,5000) [3000,8000) |
TimeWindows.ofSizeWithNoGrace(size).advanceBy(advance) |
| sliding | 시간 차(time difference) | 겹침 정도가 레코드 시각에 따라 달라집니다. 두 레코드의 타임스탬프 차가 윈도우 크기 이내면 같은 윈도우 | 레코드 타임스탬프 정렬(epoch 아님). 상·하한 모두 포함 | 집계는 SlidingWindows.ofTimeDifferenceAndGrace(...), 조인은 JoinWindows |
| session | 비활성 간격(inactivity gap) | 세션이 병합·확장됩니다 | 키마다 독립적으로 추적되고 크기도 가변입니다 | SessionWindows.ofInactivityGapWithNoGrace(gap) |
grace period — 늦게 온 데이터를 언제까지 받아 줄 것인가
grace period는 특정 윈도우에 대해 순서가 어긋난(out-of-order) 레코드를 얼마나 기다릴지를 정합니다. 공식 문서의 판정 규칙은 정확히 이렇습니다 — 레코드의 타임스탬프가 어떤 윈도우에 속하지만 현재 stream time이 그 윈도우의 끝 + grace period보다 크면, 그 레코드는 폐기되고 윈도우에 반영되지 않습니다.
윈도우 상태 저장소의 보관 기간(retention)은 grace period와 별개입니다.
Streams는 윈도우를 최소 이 기간 동안 유지하며, 기본값은 1일이고
Materialized#withRetention()으로 바꿉니다.
윈도우 계산은 새 데이터가 올 때마다 결과를 갱신해 내려보냅니다.
최종 결과만 필요하면(알림 발송, 갱신을 지원하지 않는 시스템으로 전달 등)
suppress로 윈도우가 닫힐 때까지 억제해야 합니다.
조인 매트릭스와 co-partitioning
| 피연산자 | 윈도우 | INNER | LEFT | OUTER | co-partitioning |
|---|---|---|---|---|---|
| KStream — KStream | 윈도우 필수 | 지원 | 지원 | 지원 | 필요 |
| KTable — KTable (equi-join) | 비윈도우 | 지원 | 지원 | 지원 | 필요 |
| KTable — KTable (foreign-key) | 비윈도우 | 지원 | 지원 | 미지원 | 불필요 (Streams가 내부적으로 보장) |
| KStream — KTable | 비윈도우 | 지원 | 지원 | 미지원 | 필요 |
| KStream — GlobalKTable | 비윈도우 | 지원 | 지원 | 미지원 | 불필요 |
| KTable — GlobalKTable | 해당 없음 | 미지원 | 미지원 | 미지원 | — |
co-partitioning 요건 — 공식 요건은 2개입니다
공식 문서는 "조인 시 입력 데이터가 co-partition되어 있어야 하며, 그것을 보장하는 것은 사용자의 책임이다"라고 명시합니다. 요건은 다음과 같습니다.
| 요건 | 내용 | Streams가 검증하는가 |
|---|---|---|
| 전제 · 같은 키 | equi-join은 레코드의 키를 기준으로 수행됩니다(leftRecord.key == rightRecord.key). 양쪽 입력이 키로 파티션되어 있어야 합니다 |
조인 조건 자체이므로 키가 맞지 않으면 결과가 비어 있습니다 |
| 요건 ① 같은 파티션 수 | 조인의 좌·우 입력 토픽은 파티션 수가 같아야 합니다 | 검증합니다. 파티션 배정 단계(런타임)에 다르면 TopologyException이 발생합니다 |
| 요건 ② 같은 파티셔닝 전략 | 입력 토픽에 쓰는 모든 애플리케이션이 같은 파티셔너를 써야 합니다. Java Producer면 같은 partitioner.class, Streams면 KStream#to()에 같은 StreamPartitioner |
검증할 수 없습니다. 공식 문서가 "사용자가 보장해야 한다"고 명시합니다 |
파티션 수가 다를 때
공식 문서가 권하는 절차는 이렇습니다.
- 양쪽 중 파티션 수가 적은 쪽을 찾습니다(SMALLER).
bin/kafka-topics --describe로 확인합니다. - 애플리케이션 안에서 SMALLER를
repartition으로 재분배합니다. LARGER와 같은 파티셔너를 써야 합니다.- KStream이면
KStream#repartition(Repartitioned.numberOfPartitions(...)) - KTable이면
KTable#toStream#repartition(Repartitioned.numberOfPartitions(...).toTable())
- KStream이면
- LARGER와 새 스트림/테이블을 조인합니다.
방향에 대한 문서의 권고도 함께 기억할 만합니다 — 병목을 피하려면 적은 쪽을 많은 쪽에 맞추는 것이 권장됩니다. stream-table 조인에서는 KStream을 재분배하는 쪽이 좋습니다. KTable을 재분배하면 상태 저장소가 하나 더 생길 수 있기 때문입니다. table-table 조인이면 더 작은 KTable을 재분배하는 편이 낫습니다.
// orders: 12 파티션, payments: 6 파티션 → 그대로 조인하면 TopologyException
KStream<String, Order> orders = builder.stream("orders");
KStream<String, Payment> payments = builder.stream("payments");
// 적은 쪽(payments)을 12로 맞춥니다
KStream<String, Payment> payments12 =
payments.repartition(Repartitioned.numberOfPartitions(12));
// KStream-KStream 조인은 반드시 윈도우가 필요합니다
KStream<String, Settlement> settled = orders.join(
payments12,
(order, payment) -> Settlement.of(order, payment),
JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1))
);
병렬성 — 태스크와 스레드
Streams의 병렬성 모델은 Kafka 파티션에 그대로 얹혀 있습니다.
- 각 스트림 파티션은 완전히 순서가 정해진 레코드 시퀀스이며 Kafka 토픽 파티션에 대응합니다.
- Streams는 입력 스트림 파티션에 기반해 고정된 수의 태스크를 만들고, 각 태스크에 파티션 목록을 배정합니다.
- 파티션-태스크 배정은 절대 변하지 않습니다. 태스크가 곧 고정된 병렬 단위입니다.
- 따라서 최대 병렬성은 입력 토픽의 최대 파티션 수입니다. 입력 토픽이 5파티션이면 인스턴스를 최대 5개까지 의미 있게 돌릴 수 있습니다.
- 파티션 수보다 많은 인스턴스는 기동은 되지만 유휴 상태로 남습니다. 대신 바쁜 인스턴스가 죽으면 그중 하나가 일을 이어받습니다 — 무의미한 것은 아니고 hot standby가 됩니다.
num.stream.threads와 인스턴스 수가 그 태스크를 어떻게 나눠 갖는지
| 설정 | 기본값 | 설명 | 튜닝 포인트 |
|---|---|---|---|
num.stream.threads |
1 | 인스턴스 하나가 처리에 쓰는 스레드 수. 각 스레드가 하나 이상의 태스크를 실행합니다 | 기본값이 1이므로 아무 설정도 안 하면 인스턴스당 1스레드입니다. 2.8부터 실행 중 동적으로 추가·제거할 수 있어 죽은 스레드를 재시작 없이 보충할 수 있습니다. |
num.standby.replicas |
0 | 태스크마다 유지할 standby 복제본 수 | 복구 시간과 저장 비용의 교환입니다. |
acceptable.recovery.lag |
10000 | active 태스크를 받을 만큼 "따라잡았다"고 볼 최대 lag(오프셋 수) | 리밸런스 중 처리 중단을 피하려면 이 lag이 1분 훨씬 안쪽에 복구되는 값이어야 한다고 문서가 권합니다. |
probing.rebalance.interval.ms |
600000 (10분) |
warmup 복제본이 준비됐는지 확인하는 탐색 리밸런스 주기 | 최소 1분. 배정이 균형을 이룰 때까지 반복됩니다. |
replication.factor |
-1 | changelog·리파티션 토픽의 복제 계수. -1은 브로커 기본값 사용 |
브로커의 default.replication.factor가 1이면 changelog가 복제되지 않습니다. 프로덕션에서는 3을 명시하세요. |
statestore.cache.max.bytes |
10485760 (10MiB) |
모든 스레드를 합친 상태 저장소 캐시 크기 | 캐시가 크면 중간 결과 발행이 줄어 하위 부하가 낮아지지만 결과 지연이 늘어납니다. |
topology.optimization |
none |
토폴로지 최적화 여부. all 또는 개별 최적화 목록 |
기본이 최적화 안 함입니다. StreamsBuilder#build(props)에 설정을 함께 넘겨야 적용됩니다. |
task.assignor.class |
null |
태스크 배정 구현. 기본은 HighAvailabilityTaskAssignor |
— |
exactly_once_v2
Streams의 exactly-once가 보장하는 범위는 정확히 이렇습니다 — 소스 Kafka 토픽에서 읽은 어떤 레코드든, 그 처리 결과가 출력 Kafka 토픽과 stateful 연산의 상태 저장소에 정확히 한 번 반영됩니다. 입력 토픽 오프셋 커밋, 상태 저장소 갱신, 출력 토픽 쓰기가 원자적으로 완료됩니다.
| 항목 | at_least_once (기본) | exactly_once_v2 |
|---|---|---|
commit.interval.ms | 30000 (30초) | 기본값이 100ms로 바뀝니다 |
컨슈머 isolation.level | read_uncommitted | read_committed로 설정됩니다 |
프로듀서 enable.idempotence | (프로듀서 기본값) | true로 설정됩니다 |
| 브로커 요구사항 | — | 2.5.0 이상. 기본 설정으로는 브로커 3대 이상이 필요합니다 |
인터랙티브 쿼리
Streams는 상태 저장소를 읽기 전용으로 직접 조회하는 것을 허용합니다. 저장소를 만든 애플리케이션 외부의 메서드·스레드·프로세스·애플리케이션에서도 접근할 수 있으며, 모든 저장소는 이름을 가지고 인터랙티브 쿼리는 읽기 연산만 노출합니다.
이 기능의 가치는 출력 토픽을 거치지 않고 최신 집계 결과를 바로 서빙할 수 있다는 점입니다. 집계 결과를 다시 외부 DB에 넣고 그것을 조회하는 단계가 사라집니다. 다만 상태는 인스턴스에 분산되어 있으므로, 찾는 키가 다른 인스턴스에 있으면 그 인스턴스로 요청을 라우팅해야 합니다. 그 메타데이터를 Streams가 제공합니다.
Streams Rebalance Protocol (KIP-1071)
KIP-848이 일반 컨슈머의 리밸런스 조정을 클라이언트에서 브로커로 옮겼듯이, KIP-1071은 같은 모델을 Kafka Streams 워크로드로 확장합니다. 리밸런스 이벤트마다 클라이언트가 배정을 계산하는 대신, 브로커에서 배정이 지속적으로 계산됩니다. 애플리케이션은 컨슈머 그룹이 아니라 streams group으로 등록됩니다.
| 지원됨 | 아직 지원 안 됨 |
|---|---|
|
|
streams group에는 컨슈머 그룹에 없는 NOT_READY 상태가 있습니다.
그룹 코디네이터가 토폴로지에 필요한 소스·내부 토픽이 없거나 설정이 맞지 않는 것을 감지하면
이 상태로 들어가고, 모든 멤버는 빈 배정을 받습니다.
코디네이터가 하는 검사에는 co-partition 그룹이 실제로 co-partition되어 있는지 확인하는 단계가 포함됩니다 —
앞에서 본 요건이 브로커 쪽에서도 검증되기 시작한 것입니다.
테스팅 — 브로커 없이 토폴로지를 검증합니다
CCDAK에는 Application Testing 도메인이 별도로 있고(8%),
그 안에서 가장 자주 다뤄지는 것이 TopologyTestDriver입니다.
kafka-streams-test-utils 아티팩트를 테스트 의존성으로 추가하면 씁니다.
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams-test-utils</artifactId>
<version>4.3.0</version>
<scope>test</scope>
</dependency>
TopologyTestDriver는 라이브러리 런타임을 시뮬레이션합니다 —
입력 토픽에서 레코드를 계속 가져와 토폴로지를 순회하며 처리하는 동작을 대신하고,
결과 레코드를 캡처하며 내장 상태 저장소를 조회할 수 있게 해 줍니다.
브로커가 필요하지 않습니다.
import java.time.Duration;
import java.util.Properties;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.state.KeyValueStore;
import org.junit.jupiter.api.*;
import static org.junit.jupiter.api.Assertions.*;
class WordCountTest {
private TopologyTestDriver testDriver;
private TestInputTopic<String, String> input;
private TestOutputTopic<String, Long> output;
@BeforeEach
void setUp() {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-test");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:9092"); // 실제로 접속하지 않습니다
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
testDriver = new TopologyTestDriver(WordCount.buildTopology(), props);
input = testDriver.createInputTopic("text-lines",
Serdes.String().serializer(), Serdes.String().serializer());
output = testDriver.createOutputTopic("word-counts",
Serdes.String().deserializer(), Serdes.Long().deserializer());
}
@AfterEach
void tearDown() {
// 반드시 닫아야 자원이 정리됩니다
testDriver.close();
}
@Test
void countsWordsPerKey() {
input.pipeInput("k", "kafka streams kafka");
assertEquals(new KeyValue<>("kafka", 1L), output.readKeyValue());
assertEquals(new KeyValue<>("streams", 1L), output.readKeyValue());
assertEquals(new KeyValue<>("kafka", 2L), output.readKeyValue());
assertTrue(output.isEmpty());
}
@Test
void stateStoreHoldsLatestCount() {
input.pipeInput("k", "kafka kafka kafka");
KeyValueStore<String, Long> store = testDriver.getKeyValueStore("counts-store");
assertEquals(3L, store.get("kafka"));
}
@Test
void wallClockPunctuationCanBeTriggered() {
// wall-clock time 은 드라이버가 모킹하므로 직접 전진시킵니다
testDriver.advanceWallClockTime(Duration.ofSeconds(20));
}
}
Processor API로 프로세서를 직접 작성했다면 MockProcessorContext로 단위 테스트합니다.
프로세서는 결과를 반환하지 않고 컨텍스트로 forward하므로, forward된 데이터를 캡처하는 목 컨텍스트가 필요합니다.
context.forwarded()로 검증하고, context.committed()로 커밋 호출 여부를 확인합니다.
다만 목 컨텍스트는 punctuator를 자동 실행하지 않습니다 — 캡처만 하므로 직접 호출해야 하고,
자동 발동까지 테스트하려면 TopologyTestDriver를 쓰라고 문서가 권합니다.
또한 목 컨텍스트는 changelog를 관리하지 않으므로 상태 저장소는 인메모리 + logging disabled로 등록합니다.
어떤 상황에 무엇을 쓰는가
| 목적 | 도구 | 브로커 필요 | 출처 |
|---|---|---|---|
| 프로듀서 로직 단위 테스트 | MockProducer (org.apache.kafka.clients.producer) |
불필요 | Apache Kafka kafka-clients |
| 컨슈머 로직 단위 테스트 | MockConsumer (org.apache.kafka.clients.consumer). 예외 주입은 setPollException() |
불필요 | Apache Kafka kafka-clients |
| Streams 토폴로지 검증 | TopologyTestDriver + TestInputTopic / TestOutputTopic |
불필요 | Apache Kafka kafka-streams-test-utils |
| 커스텀 Processor 단위 테스트 | MockProcessorContext |
불필요 | Apache Kafka kafka-streams-test-utils |
| 프로듀서·컨슈머 성능 측정 | kafka-producer-perf-test.sh, kafka-consumer-perf-test.sh, kafka-share-consumer-perf-test.sh |
필요 | Apache Kafka bin/ |
| 종단 지연 측정 | kafka-e2e-latency.sh |
필요 | Apache Kafka bin/ |
| 실제 브로커 대상 통합 테스트 | Testcontainers, spring-kafka의 @EmbeddedKafka 등 |
필요 (컨테이너/임베디드) |
외부 프로젝트 — 버전과 API는 각 프로젝트 문서로 확인하세요 |
| 목 데이터 생성 | kafka-connect-datagen 커넥터, ksqlDB의 datagen |
필요 | Confluent 제공 — Apache Kafka 배포판에 없습니다 |
ksqlDB — Streams 위의 SQL 계층
ksqlDB의 가치는 Java 애플리케이션을 배포하지 않고 SQL 문장으로 스트림 처리를 정의한다는 데 있습니다.
Streams의 KStream이 STREAM, KTable이 TABLE에 대응합니다.
-- 기존 토픽을 STREAM 으로 선언 (KStream 에 대응: 각 레코드가 독립 사건)
CREATE STREAM orders (
order_id BIGINT KEY,
user_id VARCHAR,
amount DOUBLE
) WITH (
KAFKA_TOPIC = 'orders',
VALUE_FORMAT = 'AVRO'
);
-- 기존 토픽을 TABLE 로 선언 (KTable 에 대응: 키별 최신 상태)
CREATE TABLE users (
user_id VARCHAR PRIMARY KEY,
country VARCHAR,
tier VARCHAR
) WITH (
KAFKA_TOPIC = 'users',
VALUE_FORMAT = 'AVRO'
);
-- 지속 질의(persistent query): 결과가 새 토픽으로 계속 흘러갑니다
CREATE STREAM high_value_orders WITH (
KAFKA_TOPIC = 'high_value_orders',
VALUE_FORMAT = 'AVRO'
) AS
SELECT o.order_id, o.amount, u.country, u.tier
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.amount > 1000
EMIT CHANGES;
| 기준 | Kafka Streams | ksqlDB |
|---|---|---|
| 개발 언어 | Java / Scala | SQL |
| 배포 형태 | 내 애플리케이션 안에 임베드 | ksqlDB 서버 클러스터에 질의를 등록 |
| 표현력 | 임의 로직 가능 (Processor API 포함) | SQL로 표현되는 범위. UDF로 확장 |
| 운영 주체 | 애플리케이션 팀 | ksqlDB 서버를 운영하는 팀 |
| 언제 쓰나 | 기존 서비스에 스트림 처리를 넣을 때, 복잡한 로직·외부 연동이 필요할 때 | 탐색적 분석, 빠른 프로토타이핑, SQL에 익숙한 팀이 파이프라인을 직접 만들 때 |
| Apache Kafka 포함 | 포함 | 미포함 (Confluent) |
무엇을 모니터링해야 하는가
| 메트릭 | MBean 그룹 | 무엇을 알려 주는가 |
|---|---|---|
failed-stream-threads | kafka.streams:type=stream-metrics | 클라이언트 시작 이후 실패한 스트림 스레드 수. 0이 아니면 처리 용량이 줄어든 상태입니다 |
alive-stream-threads | kafka.streams:type=stream-metrics | 현재 살아 있는(또는 리밸런스 참여 중인) 스레드 수 |
client-state / thread-state | stream-metrics / stream-thread-metrics | 클라이언트·스레드 상태의 enum ordinal 값 |
process-rate, process-latency-avg | kafka.streams:type=stream-thread-metrics | 초당 처리 레코드 수와 평균 처리 시간 |
commit-latency-avg | kafka.streams:type=stream-thread-metrics | 이 스레드의 모든 태스크에 대한 평균 커밋 실행 시간(ms) |
poll-latency-avg | kafka.streams:type=stream-thread-metrics | 컨슈머 폴링 평균 시간 |
blocked-time-ns-total | kafka.streams:type=stream-thread-metrics | 브로커 때문에 블록된 총 시간. 병목이 내 코드인지 브로커인지 가르는 지표입니다 |
task-created-rate | kafka.streams:type=stream-thread-metrics | 초당 태스크 생성률. 계속 높으면 리밸런스 스톰입니다 |
record-e2e-latency-max | stream-processor-node-metrics / stream-state-metrics | 레코드 타임스탬프와 처리 시점 시스템 시각을 비교한 종단 지연. recording level이 info여서 프로덕션에서도 켜 둘 수 있습니다 |
restore-rate | kafka.streams:type=stream-state-metrics | 상태 저장소 복원 속도. 재시작 후 얼마나 남았는지 추정에 씁니다 |
흔한 오해
시험 포인트 정리
확인 문제
네 가지 문항 유형(단일 선택 · 복수 선택 · 연결형 · 순서 배열)이 섞여 있습니다. 리파티션 유발 여부와 co-partitioning 요건은 반드시 맞히고 넘어가세요.
이어서 볼 곳
- 11장 · 운영 기초 보안·모니터링·성능 튜닝과 Share Groups.
- 예제 9 · Streams 실시간 집계 윈도우 집계와 상태 저장소를 실제로 돌려 봅니다.
- Streams 치트시트 DSL 연산자 표(리파티션 유발 여부), 윈도우·조인 매트릭스, 설정 15개.
- CCDAK · Kafka Streams 도메인 12% 도메인의 압축 정리와 자주 나오는 함정.
- CCDAK · Application Testing 도메인 8% 도메인. TopologyTestDriver와 목 객체 정리.
- 9장 · Kafka Connect SMT의 한계와 Streams로 넘어가야 하는 경계.
공식 문서 출처
이 장의 설정 기본값·API·의미론은 모두 아래에서 확인했습니다 (Apache Kafka 4.3 문서 기준).
- Streams Core Concepts — 토폴로지·source/sink processor, event/processing/ingestion time, stream time, 출력 타임스탬프 규칙, stream-table 이원성, 윈도우와 grace period, exactly-once v2, 순서 어긋남 처리
- Streams Architecture — 파티션·태스크 관계, 최대 병렬성,
StreamsPartitionAssignor교체 불가, changelog와 로그 컴팩션, standby replica, 랙 인식 배정 - Streams DSL — KStream/KTable/GlobalKTable 정의와
("alice",1)→("alice",3)예시, stateless/stateful 연산 표와 리파티션 표시 규칙, 조인 매트릭스, co-partitioning 요구사항, 윈도우 4종 - Streams Configs —
num.stream.threads1,num.standby.replicas0,replication.factor-1,commit.interval.ms30000(EOS 시 100),statestore.cache.max.bytes,max.task.idle.ms,topology.optimization - Testing a Streams Application —
kafka-streams-test-utils,TopologyTestDriver,TestInputTopic/TestOutputTopic,advanceWallClockTime,MockProcessorContext의 한계 - Streams Rebalance Protocol — KIP-1071, 지원/미지원 목록, 그룹 레벨 설정,
NOT_READY상태, 오프라인 마이그레이션 절차 - Interactive Queries — 읽기 전용 상태 저장소 조회
- Upgrading Apache Kafka — 4.1 Streams 리밸런스 프로토콜 Early Access, 4.2 production-ready, KAFKA-20254와 4.2.1 수정
- Kafka Streams Monitoring —
failed-stream-threads,blocked-time-ns-total,record-e2e-latency-*의 recording level - Managing Consumer Groups —
kafka-consumer-groups.sh --describe출력,kafka-groups.sh - kafka-streams-test-utils (Apache Kafka 4.3 소스) —
TopologyTestDriver·TestInputTopic·MockProcessorContext의 실제 패키지 경로