학습 목표

클러스터가 아니라 라이브러리입니다

스트림 처리 프레임워크는 보통 자체 클러스터와 리소스 매니저를 요구합니다. Kafka Streams는 그 방향을 택하지 않았습니다. 공식 문서의 표현대로 "어떤 Java 애플리케이션에도 쉽게 임베드할 수 있는 단순하고 가벼운 클라이언트 라이브러리"이고, 내부 메시징 계층으로 Apache Kafka 자체 외에 어떤 외부 시스템에도 의존하지 않습니다. 문서는 "Kafka Streams는 리소스 매니저가 아니며, 스트림 처리 애플리케이션이 도는 곳이면 어디서든 돈다"고 명시합니다.

Kafka Streams가 Consumer/Producer API 직접 사용과 다른 지점
관점 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)로 계산 로직을 정의합니다. 토폴로지는 스트림 프로세서(노드)를 스트림(엣지)으로 연결한 그래프이며, 두 개의 특별한 프로세서가 있습니다.

토폴로지를 정의하는 방법은 두 가지입니다 — Kafka Streams DSL(map, filter, join, 집계 등을 바로 제공)과 더 저수준인 Processor API(커스텀 프로세서를 직접 정의하고 상태 저장소에 접근). 토폴로지는 논리적 추상일 뿐이고, 런타임에는 병렬 처리를 위해 인스턴스화되고 복제됩니다.

Streams 토폴로지 — sub-topology 경계와 리파티션 토픽 Kafka Streams 토폴로지가 리파티션 토픽에서 두 개의 sub-topology 로 갈리는 모습을 그린 그림입니다. 첫 번째 sub-topology 는 orders 토픽을 읽는 source 노드에서 시작해 filter, selectKey 를 지나 리파티션 토픽으로 쓰는 sink 노드로 끝납니다. selectKey 가 키를 바꾸므로 스트림에 리파티션 표시가 붙고, 뒤에서 집계를 하면 실제로 리파티션이 일어납니다. 가운데의 리파티션 토픽이 두 sub-topology 를 잇는 경계입니다. 두 번째 sub-topology 는 그 리파티션 토픽을 읽는 source 노드에서 시작해 상태 저장소를 쓰는 집계 노드, KTable 을 KStream 으로 바꾸는 노드를 지나 결과 토픽으로 쓰는 sink 노드로 끝납니다. 태스크는 sub-topology 단위로 만들어지므로 경계가 병렬성의 단위도 함께 나눕니다. 리파티션 토픽 이름은 application.id 와 연산자 이름과 repartition 접미사로 이루어지고, cleanup.policy 는 delete, retention.ms 는 무한이며 처리된 데이터는 Streams 가 자동으로 지웁니다. 토폴로지 — sub-topology 는 리파티션 토픽에서 갈립니다 sub-topology 0 source orders filter stateless selectKey 리파티션 표시됨 sink → 리파티션 토픽 리파티션 토픽 (내부) sub-topology 경계 sub-topology 1 source 리파티션 토픽 count 상태 저장소 사용 toStream KTable → KStream sink order-counts 태스크는 sub-topology 단위로 만들어집니다 — 경계가 병렬성의 단위도 나눕니다. 리파티션 토픽 이름은 {application.id}-{연산자}-repartition 이고 cleanup.policy=delete 입니다.
토폴로지 구조 — source processor에서 시작해 여러 processor를 거쳐 sink processor로 끝나는 그래프와, 리파티션 토픽이 sub-topology를 나누는 경계
WordCount — DSL로 정의한 최소 토폴로지
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)이 들어올 때 사용자별 값을 합산하면,

세 추상의 의미론 비교 (Apache Kafka 4.3 Streams DSL 문서 기준)
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)
비용 로컬 저장 공간과 브로커 부하가 큼 (토픽 전체를 각 인스턴스가 읽음)
어떤 데이터에 쓰나 클릭, 결제 트랜잭션, 센서 측정, 로그 사용자 프로필, 재고 수량, 계정 잔액 국가 코드, 환율처럼 작고 자주 변하지 않는 참조 데이터
KStream · KTable · GlobalKTable — 같은 입력을 다르게 해석합니다 같은 입력 레코드 다섯 건을 세 가지 추상화가 어떻게 해석하는지 비교한 그림입니다. 입력은 시간 순서대로 키 A 에 값 1, 키 B 에 값 2, 키 A 에 값 3, 키 B 에 값 null, 키 A 에 값 4 입니다. KStream 은 레코드 스트림이므로 다섯 건을 각각 독립적인 이벤트로 다룹니다. 같은 키가 반복되어도 덮어쓰지 않고 값이 null 인 것도 하나의 이벤트로 전달됩니다. KTable 은 변경로그로 해석하므로 키별로 마지막 값만 남는 upsert 가 됩니다. 결과는 A 가 4 이고, B 는 값이 null 인 레코드를 tombstone 으로 보아 삭제됩니다. GlobalKTable 은 해석 자체는 KTable 과 같지만 저장 방식이 다릅니다. 입력 토픽의 모든 파티션을 모든 애플리케이션 인스턴스가 전부 복제해 가지고 있습니다. 그래서 조인할 때 co-partitioning 이 필요하지 않고, 키가 아닌 값으로도 조회할 수 있습니다. 대신 인스턴스마다 전체 데이터를 들고 있어야 하므로 작은 참조 데이터에 적합합니다. KStream vs KTable vs GlobalKTable — 입력은 같고 해석이 다릅니다 입력 토픽 t1 t2 t3 t4 t5 key A value 1 key B value 2 key A value 3 key B value null key A value 4 시간 → KStream 레코드 스트림 · 각 이벤트 독립 A = 1 B = 2 A = 3 B = null ← 이것도 이벤트 A = 4 KTable 변경로그 · 키별 upsert 최종 상태 A = 4 1·3 은 덮어써짐 B 삭제됨 (tombstone) 키 2개 입력 → 행 1개 남음 상태 저장소 + changelog 사용 GlobalKTable 해석은 KTable 과 동일 최종 상태 (KTable 과 같음) A = 4 B 삭제됨 (tombstone) 다른 점 — 모든 파티션을 모든 인스턴스가 전부 복제 GlobalKTable 은 전체 복제라서 조인 시 co-partitioning 이 필요 없고 값으로도 조회할 수 있습니다 (D-093). 대신 인스턴스마다 전체 데이터를 들고 있어야 하므로 작은 참조 데이터에만 씁니다. 같은 토픽을 builder.stream() 으로 읽으면 KStream, builder.table() 로 읽으면 KTable 입니다. 데이터가 다른 것이 아니라 해석이 다른 것입니다.
같은 입력에 대한 세 가지 해석 — ("alice",1) → ("alice",3)을 KStream(합 4) · KTable(합 3) · GlobalKTable(전 파티션 복제)로 읽었을 때의 결과 차이

stateless와 stateful, 그리고 리파티션

DSL 연산을 분류하는 축은 두 개이고 서로 독립입니다 — 상태 저장소를 만드는가리파티션을 유발하는가입니다. map은 stateless이지만 리파티션을 유발하고, count는 stateful이지만 그 자체가 리파티션을 유발하지는 않습니다. 이 두 축을 섞으면 문제를 틀립니다.

주요 DSL 연산의 상태 여부와 리파티션 유발 여부 (Apache Kafka 4.3 DSL 문서 기준)
연산 상태 리파티션 비고
filter / filterNotstateless유발 안 함키를 바꾸지 않습니다
mapValuesstateless유발 안 함값만 바꿉니다
flatMapValuesstateless유발 안 함값만 바꾸며 개수를 늘립니다
mapstateless표시함키를 바꿀 수 있으므로 리파티션 대상으로 표시(mark)됩니다
flatMapstateless표시함가능하면 flatMapValues를 쓰라고 문서가 권합니다
selectKeystateless표시함키를 바꾸는 것이 목적입니다
groupByKeystateless
(그룹화만)
표시됐을 때만스트림이 리파티션 대상으로 표시된 경우에만 실제로 리파티션합니다
groupBystateless
(그룹화만)
항상키를 새로 정하므로 언제나 리파티션합니다. 가능하면 groupByKey를 쓰세요
repartitionstateless항상파티션 수를 지정해 명시적으로 재분배. 내부 토픽은 Streams가 관리·자동 삭제합니다
branch / merge / peek / foreach / printstateless유발 안 함foreach·print는 종단 연산입니다
count / reduce / aggregatestateful앞 연산에 달림상태 저장소 + changelog 토픽이 생깁니다
모든 윈도우 연산stateful앞 연산에 달림윈도우 상태 저장소를 씁니다
모든 조인stateful표시됐을 때만양쪽이 표시됐으면 양쪽이 리파티션됩니다
suppressstateful유발 안 함윈도우 최종 결과만 내보낼 때 씁니다
cogroupstateful유발 안 함입력이 이미 그룹화되어 있는 것이 전제이므로 리파티션을 유발하지 않습니다
stateless · stateful · 리파티션 유발 연산 — map 과 mapValues 의 차이 Kafka Streams DSL 연산을 세 갈래로 분류한 그림입니다. 첫째 열 stateless 는 상태 저장소가 필요 없는 연산입니다. filter, filterNot, map, mapValues, flatMap, flatMapValues, selectKey, split, merge, peek, foreach, print, to, groupByKey, groupBy 가 여기에 속합니다. branch 와 through 는 4.0 에서 제거되었으므로 각각 split 과 repartition 을 씁니다. 둘째 열 stateful 은 상태 저장소가 필요한 연산입니다. aggregate, reduce, count, 윈도우 집계, KStream-KStream 조인, KStream-KTable 조인, KTable-KTable 조인, suppress, 상태 저장소를 연결한 process 가 여기에 속합니다. 셋째 열은 리파티션을 유발하는지에 따른 분류입니다. groupBy 와 repartition 은 항상 리파티션을 유발합니다. map, flatMap, selectKey 는 키를 바꿀 수 있으므로 스트림에 리파티션 표시를 남기고, 그 뒤에 그룹화나 조인이 오면 실제로 리파티션이 일어납니다. mapValues, flatMapValues, filter, peek 은 키를 건드리지 않으므로 리파티션을 유발하지 않습니다. groupByKey 는 앞에서 표시가 붙지 않았다면 리파티션을 하지 않습니다. 핵심 대비는 map 과 mapValues 입니다. map 은 키를 바꿀 수 있으므로 리파티션 토픽을 거쳐 네트워크 왕복이 생기고, mapValues 는 값만 바꾸므로 그대로 로컬에서 처리됩니다. stateless · stateful · 리파티션 유발 — 세 축으로 나눠 외웁니다 stateless — 저장소 없음 filter · filterNot map · mapValues flatMap · flatMapValues selectKey split merge peek · foreach · print to groupByKey · groupBy toStream · repartition 그룹화는 stateless 입니다 — 상태는 그 뒤 집계에서 생깁니다. stateful — 저장소 필요 aggregate reduce count windowedBy(…) + 집계 join / leftJoin / outerJoin suppress process (저장소 연결 시) 상태 저장소에는 changelog 토픽이 따라옵니다 (D-095). KStream-GlobalKTable 조인은 전체 복제본을 조회하는 형태라 리파티션이 필요 없습니다. 리파티션 유발 여부 항상 유발 groupBy repartition 키를 바꿔 “표시”만 남김 map flatMap selectKey → 뒤에 그룹화·조인이 오면 발생 유발하지 않음 mapValues · flatMapValues filter · peek · groupByKey map 키를 바꿀 수 있음 → 리파티션 유발 리파티션 토픽에 쓰고 다시 읽습니다. 네트워크·디스크 왕복 + sub-topology 분리. mapValues 값만 바꿈 → 유발 안 함 키가 그대로이므로 로컬에서 끝납니다. 값만 바꿀 때는 항상 mapValues 를 쓰세요. Streams 는 키가 실제로 바뀌었는지 보지 않고 키를 바꿀 수 있는 연산인지로 판단합니다. 그래서 키를 그대로 반환하는 map 에도 리파티션 표시가 붙습니다.
stateless · stateful · 리파티션 유발의 3열 분류 — map은 유발하고 mapValues는 유발하지 않는다는 대비가 중심

상태 저장소와 changelog

stateful 연산을 쓰면 DSL이 상태 저장소(state store)를 자동으로 만들고 관리합니다. 기본 구현은 디스크 기반 키-값 저장소(RocksDB)이고, 인메모리 해시맵도 선택할 수 있습니다. 내고장성은 changelog 토픽이 담당합니다.

문제는 복원 시간입니다. 태스크 재초기화 비용은 대부분 changelog 재생 시간이 차지합니다. 이를 줄이는 장치가 standby replica입니다. num.standby.replicas를 올리면 상태의 완전 복제본을 다른 인스턴스에 미리 유지하고, 태스크가 이동할 때 이미 최신 상태를 가진 인스턴스로 배정합니다. 2.6부터는 완전히 따라잡은 로컬 복사본을 가진 인스턴스가 존재한다면 반드시 그 인스턴스에 배정됩니다.

상태 저장소와 changelog — RocksDB · changelog 토픽 · standby 복구 Kafka Streams 의 상태 저장이 어떻게 이루어지고 장애에서 어떻게 복구되는지 그린 그림입니다. 위쪽은 쓰기 경로입니다. 스트림 스레드가 레코드를 처리하면서 로컬 RocksDB 저장소에 값을 쓰고, 같은 변경을 changelog 토픽에도 씁니다. changelog 는 로컬 상태의 백업이며 태스크마다 자기 파티션을 하나 가집니다. key-value 저장소의 changelog 는 cleanup.policy 가 compact 이고, 윈도우 저장소의 changelog 는 delete 와 compact 를 함께 씁니다. 아래쪽은 복구 경로입니다. standby 레플리카가 없으면 태스크가 다른 인스턴스로 옮겨간 뒤 changelog 를 처음부터 재생해야 하므로 상태가 클수록 복구가 오래 걸립니다. num.standby.replicas 를 1 이상으로 두면 다른 인스턴스가 미리 changelog 를 따라 읽어 거의 최신 상태를 들고 있으므로 훨씬 빨리 처리를 재개합니다. num.standby.replicas 의 기본값은 0 입니다. changelog 토픽 이름은 application.id 와 저장소 이름과 changelog 접미사로 이루어집니다. 상태 저장소와 changelog — 로컬 상태는 항상 Kafka 에 백업됩니다 쓰기 경로 — 처리할 때마다 두 곳에 씁니다 스트림 스레드 task 0_0 count · aggregate 처리 로컬 상태 저장소 RocksDB (디스크) 태스크 전용 · 공유 안 함 changelog 토픽 (내부) {application.id}-{store}-changelog 태스크마다 자기 파티션 하나 cleanup.policy=compact 인스턴스가 죽었을 때 — standby 가 있는지로 복구 시간이 갈립니다 num.standby.replicas = 0 기본값 태스크가 다른 인스턴스로 옮겨간 뒤 changelog 를 처음부터 재생해 상태를 만듭니다. 상태가 크면 복구가 오래 걸립니다. 그동안 그 파티션은 처리되지 않습니다. num.standby.replicas ≥ 1 다른 인스턴스가 changelog 를 미리 따라 읽어 거의 최신 상태를 들고 있습니다. 거의 즉시 처리를 재개합니다. 대신 디스크·네트워크를 더 씁니다. 알아 둘 것 changelog 는 compact 라서 키별 최신값만 남습니다 — 그래서 재생하면 상태가 그대로 복원됩니다. 윈도우 저장소의 changelog 는 cleanup.policy=delete,compact 이고 보관 기간은 윈도우 크기 + 24시간입니다. changelog 토픽을 사람이 지우면 상태 복구가 불가능해집니다. 내부 토픽은 손대지 마세요.
상태 저장소와 changelog — 태스크의 로컬 RocksDB, 그에 대응하는 changelog 파티션, 그리고 standby replica가 미리 따라잡아 두어 복구 시간을 줄이는 구조

시간 — event time · processing time · ingestion time

세 가지 시간 개념 (Apache Kafka 4.3 Streams 핵심 개념 문서 기준)
개념언제의 시각인가누가 부여하는가
event time 이벤트가 소스에서 실제로 발생한 시점 이벤트를 만든 쪽 (예: GPS 센서가 위치 변화를 감지한 순간)
processing time 스트림 처리 애플리케이션이 그 레코드를 처리한 시점 처리 애플리케이션. event time보다 밀리초~수 시간 늦을 수 있습니다
ingestion time 브로커가 토픽 파티션에 저장한 시점 브로커. 레코드가 끝까지 처리되지 않아도 ingestion time은 존재합니다

Streams는 TimestampExtractor 인터페이스로 모든 레코드에 타임스탬프를 부여합니다. 이 값이 진행 상황을 나타내는 stream time이고, 실제 실행 시각인 wall-clock time과 구분됩니다. stream time은 새 레코드가 프로세서에 도착할 때만 전진합니다 — 데이터가 멈추면 시간도 멈춥니다. 윈도우가 닫히지 않는 현상의 원인이 대부분 이것입니다.

타임스탬프 관련 Streams 설정 (Apache Kafka 4.3 기본값)
설정기본값설명튜닝 포인트
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로 올라갑니다.

출력 레코드의 타임스탬프는 어떻게 정해지는가

출력 타임스탬프 규칙 (Apache Kafka 4.3 Streams 핵심 개념 문서 기준)
상황출력 타임스탬프
입력 레코드를 처리해 출력 (context.forward() 등)입력 레코드의 타임스탬프를 그대로 상속
주기 함수(Punctuator#punctuate())에서 출력스트림 태스크의 현재 내부 시각
stream-stream / table-table 조인max(left.ts, right.ts)
stream-table 조인stream 쪽 레코드의 타임스탬프
집계키별(또는 윈도우별)로 기여한 모든 입력 중 최대 타임스탬프
stateless 연산입력 타임스탬프 통과. flatMap류는 모든 출력이 같은 값을 상속

윈도우 4종

윈도우는 stateful 연산을 위해 같은 키를 가진 레코드를 시간 구간으로 다시 묶는 장치입니다. 공식 문서가 강조하는 전제: 윈도우는 레코드 키별로 따로 추적됩니다.

윈도우 4종 비교 (Apache Kafka 4.3 DSL 문서 기준)
종류 정의 요소 겹침 정렬 기준 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)
윈도우 4종 비교 — 같은 이벤트 열이 어떻게 묶이는가 같은 키의 이벤트가 6초, 8초, 9초, 14초에 도착했을 때 tumbling, hopping, sliding, session 네 가지 윈도우가 각각 어떻게 묶는지 비교한 그림입니다. tumbling 은 크기 5초이고 에포크에 정렬되므로 0에서 5, 5에서 10, 10에서 15, 15에서 20 구간이 겹치지 않게 이어집니다. 6, 8, 9 는 5에서 10 구간에, 14 는 10에서 15 구간에 들어가고 각 이벤트는 정확히 한 윈도우에만 속합니다. hopping 은 크기 5초에 전진 간격 3초이므로 윈도우가 겹칩니다. 3에서 8 은 6 을, 6에서 11 은 6과 8과 9 를, 9에서 14 는 9 를, 12에서 17 은 14 를 담습니다. 한 이벤트가 여러 윈도우에 들어갑니다. sliding 은 시간 차 5초 기준이며 윈도우 경계가 에포크가 아니라 레코드 타임스탬프에 정렬되고 양 끝이 모두 포함됩니다. 레코드가 윈도우에 들어올 때와 빠져나갈 때마다 새 윈도우가 생겨 1에서 6, 3에서 8, 4에서 9, 7에서 12, 9에서 14, 10에서 15 여섯 개가 만들어집니다. session 은 비활동 간격 3초 기준입니다. 6, 8, 9 는 간격이 3 이하라 하나의 세션 6에서 9 로 합쳐지고, 9와 14 사이 간격이 5로 3보다 크므로 14 는 자기 혼자 세션이 됩니다. session 윈도우는 크기가 고정되어 있지 않고 데이터가 크기를 정합니다. 윈도우 4종 — 같은 이벤트 열, 다른 묶음 이벤트 (같은 키) 도착 시각 = 6s · 8s · 9s · 14s · 시간 윈도우 크기 5s · hopping 전진 3s · session 비활동 간격 3s 이벤트 6 8 9 14 0s 5s 10s 15s 20s ① tumbling — 고정 크기 · 겹치지 않음 · 에포크 정렬 · 이벤트는 한 윈도우에만 [0,5) 없음 [5,10) → 6·8·9 [10,15) → 14 [15,20) 없음 ② hopping — 고정 크기 · 3초마다 전진 · 겹침 · 한 이벤트가 여러 윈도우에 [3,8) → 6 [6,11) → 6·8·9 [9,14) → 9 [12,17) → 14 ③ sliding — 크기 고정 · 경계가 레코드 타임스탬프에 정렬 · 양 끝 포함 [1,6] → 6 [3,8] → 6·8 [4,9] → 6·8·9 [7,12] → 8·9 [9,14] → 9·14 [10,15] → 14 ④ session — 크기 가변 · 비활동 간격이 기준 · 데이터가 크기를 정함 [6,9] 6·8·9 병합 [14,14] 단독 세션 tumbling 은 hopping 의 특수한 경우입니다 — 크기 = 전진 간격 이면 tumbling 이 됩니다. session — 간격이 3 이하면 같은 세션으로 병합됩니다. 9 에서 14 는 간격 5 로 3보다 크므로 새 세션이 열립니다. 시간 윈도우 [a,b) 는 시작 포함·끝 제외, sliding [a,b] 는 양쪽 포함입니다. session 윈도우는 키마다 시작·끝이 다르고 크기도 서로 다릅니다.
윈도우 4종 비교 — 같은 이벤트 열에 tumbling · hopping · sliding · session 윈도우를 적용했을 때 각 레코드가 어느 윈도우에 속하는지

grace period — 늦게 온 데이터를 언제까지 받아 줄 것인가

grace period는 특정 윈도우에 대해 순서가 어긋난(out-of-order) 레코드를 얼마나 기다릴지를 정합니다. 공식 문서의 판정 규칙은 정확히 이렇습니다 — 레코드의 타임스탬프가 어떤 윈도우에 속하지만 현재 stream time이 그 윈도우의 끝 + grace period보다 크면, 그 레코드는 폐기되고 윈도우에 반영되지 않습니다.

윈도우 상태 저장소의 보관 기간(retention)은 grace period와 별개입니다. Streams는 윈도우를 최소 이 기간 동안 유지하며, 기본값은 1일이고 Materialized#withRetention()으로 바꿉니다.

윈도우 계산은 새 데이터가 올 때마다 결과를 갱신해 내려보냅니다. 최종 결과만 필요하면(알림 발송, 갱신을 지원하지 않는 시스템으로 전달 등) suppress로 윈도우가 닫힐 때까지 억제해야 합니다.

조인 매트릭스와 co-partitioning

지원되는 조인 조합 (Apache Kafka 4.3 DSL 문서 기준)
피연산자 윈도우 INNER LEFT OUTER co-partitioning
KStream — KStream 윈도우 필수 지원 지원 지원 필요
KTable — KTable (equi-join) 비윈도우 지원 지원 지원 필요
KTable — KTable (foreign-key) 비윈도우 지원 지원 미지원 불필요 (Streams가 내부적으로 보장)
KStream — KTable 비윈도우 지원 지원 미지원 필요
KStream — GlobalKTable 비윈도우 지원 지원 미지원 불필요
KTable — GlobalKTable 해당 없음 미지원 미지원 미지원
조인 매트릭스와 co-partitioning 요건 Kafka Streams 의 조인 조합별로 윈도우 필요 여부, INNER · LEFT · OUTER 지원 여부, co-partitioning 필요 여부를 정리한 표입니다. KStream-KStream 조인은 윈도우가 필수이며 INNER, LEFT, OUTER 를 모두 지원하고 co-partitioning 이 필요합니다. KStream-KTable 조인은 윈도우가 없고 INNER 와 LEFT 만 지원하며 co-partitioning 이 필요합니다. KStream-GlobalKTable 조인은 윈도우가 없고 INNER 와 LEFT 만 지원하며 co-partitioning 이 필요하지 않습니다. KTable-KTable 조인은 윈도우가 없고 INNER, LEFT, OUTER 를 모두 지원하며 co-partitioning 이 필요합니다. KTable-KTable 외래키 조인은 윈도우가 없고 INNER 와 LEFT 만 지원하며 co-partitioning 이 필요하지 않습니다. 공식 문서가 co-partitioning 의 요건으로 명시하는 것은 두 가지입니다. 첫째 양쪽 입력 토픽의 파티션 수가 같아야 합니다. 둘째 그 토픽에 쓰는 모든 애플리케이션의 파티셔너가 같아야 합니다. equi-join 이 키로 매칭하므로 양쪽이 같은 키로 파티셔닝되어 있어야 한다는 것은 그 전제입니다. Kafka Streams 는 이 중 파티션 수만 실행 중에 검증하며 다르면 예외를 던집니다. 파티셔너가 같은지는 검증할 수 없으므로 사용자가 보장해야 합니다. GlobalKTable 조인은 모든 인스턴스가 전체 데이터를 가지고 있어서, 외래키 조인은 Streams 가 내부적으로 co-partitioning 을 맞춰 주기 때문에 요건이 면제됩니다. 조인 매트릭스 — 윈도우가 필요한 조인은 하나뿐입니다 조인 조합 윈도우 INNER LEFT OUTER co-partitioning KStream - KStream 필수 지원 지원 지원 필요 KStream - KTable 없음 지원 지원 불가 필요 KStream - GlobalKTable 없음 지원 지원 불가 불필요 KTable - KTable 없음 지원 지원 지원 필요 KTable - KTable (외래키) 없음 지원 지원 불가 불필요 co-partitioning 요건 — 공식 문서가 요건으로 명시하는 것은 2가지입니다 1 파티션 수 동일 — 다르면 실행 중 예외가 발생합니다 (Streams 가 검증하는 유일한 항목). 2 파티셔너 동일 — 쓰는 쪽 파티셔너가 같아야 합니다. Streams 는 검증할 수 없습니다. 전제 equi-join 은 키로 매칭하므로 양쪽이 같은 키로 파티셔닝되어 있어야 합니다. co-partitioning 이 면제되는 두 경우 GlobalKTable 조인 — 모든 인스턴스가 전체 파티션을 다 가지고 있습니다. KTable 외래키 조인 — Streams 가 내부에서 co-partitioning 을 맞춰 줍니다. 파티션 수가 다르면 작은 쪽을 큰 쪽에 맞춰 리파티션하는 것이 권장됩니다. 스트림-테이블 조인이면 KStream 쪽을, 테이블-테이블이면 작은 KTable 쪽을 리파티션하세요.
조인 매트릭스와 co-partitioning — 4가지 조인 조합의 윈도우 필요 여부와, 파티션 수 동일 · 파티셔너 동일이라는 co-partitioning 2요건(키 동일은 equi-join의 전제)

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

공식 문서는 "조인 시 입력 데이터가 co-partition되어 있어야 하며, 그것을 보장하는 것은 사용자의 책임이다"라고 명시합니다. 요건은 다음과 같습니다.

co-partitioning 요구사항
요건내용Streams가 검증하는가
전제 · 같은 키 equi-join은 레코드의 키를 기준으로 수행됩니다(leftRecord.key == rightRecord.key). 양쪽 입력이 키로 파티션되어 있어야 합니다 조인 조건 자체이므로 키가 맞지 않으면 결과가 비어 있습니다
요건 ① 같은 파티션 수 조인의 좌·우 입력 토픽은 파티션 수가 같아야 합니다 검증합니다. 파티션 배정 단계(런타임)에 다르면 TopologyException이 발생합니다
요건 ② 같은 파티셔닝 전략 입력 토픽에 쓰는 모든 애플리케이션이 같은 파티셔너를 써야 합니다. Java Producer면 같은 partitioner.class, Streams면 KStream#to()에 같은 StreamPartitioner 검증할 수 없습니다. 공식 문서가 "사용자가 보장해야 한다"고 명시합니다

파티션 수가 다를 때

공식 문서가 권하는 절차는 이렇습니다.

  1. 양쪽 중 파티션 수가 적은 쪽을 찾습니다(SMALLER). bin/kafka-topics --describe로 확인합니다.
  2. 애플리케이션 안에서 SMALLER를 repartition으로 재분배합니다. LARGER와 같은 파티셔너를 써야 합니다.
    • KStream이면 KStream#repartition(Repartitioned.numberOfPartitions(...))
    • KTable이면 KTable#toStream#repartition(Repartitioned.numberOfPartitions(...).toTable())
  3. 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 파티션에 그대로 얹혀 있습니다.

태스크와 스레드 병렬성 — 파티션 수가 병렬성의 상한입니다 입력 토픽의 파티션 여섯 개가 태스크 여섯 개로 이어지고, 그 태스크가 애플리케이션 인스턴스 두 대의 스트림 스레드 세 개씩에 하나씩 배정된 모습을 그린 그림입니다. 태스크 수는 sub-topology 가 읽는 입력 토픽의 파티션 수로 정해지며 실행 중에 바뀌지 않습니다. num.stream.threads 는 인스턴스 하나가 몇 개의 태스크를 동시에 실행할지를 정하는 설정이고 기본값은 1 입니다. 이 그림에서는 num.stream.threads 를 3 으로 두고 인스턴스를 두 대 띄웠으므로 스레드 총합이 여섯 개가 되어 태스크 여섯 개와 정확히 맞습니다. 스레드 총합이 태스크 수보다 많으면 남는 스레드는 아무 태스크도 받지 못하고 유휴 상태가 됩니다. 인스턴스를 더 띄워도 파티션 수를 넘는 병렬성은 얻을 수 없습니다. 처리량을 더 늘리려면 입력 토픽의 파티션 수를 늘려야 합니다. 스레드 사이에는 공유 상태가 없으므로 스레드 간 조정도 필요하지 않습니다. 태스크와 스레드 — 파티션 수 = 태스크 수 = 병렬성 상한 입력 토픽 파티션 6 P0 P1 P2 P3 P4 P5 태스크 6개 고정 0_0 0_1 0_2 0_3 0_4 0_5 인스턴스 A — num.stream.threads=3 thread-1 0_0 태스크 1개 thread-2 0_1 태스크 1개 thread-3 0_2 태스크 1개 인스턴스 B — num.stream.threads=3 thread-1 0_3 태스크 1개 thread-2 0_4 태스크 1개 thread-3 0_5 태스크 1개 규칙 태스크 수 = sub-topology 가 읽는 입력 토픽의 파티션 수. 실행 중에는 바뀌지 않습니다. num.stream.threads (기본 1) 는 인스턴스 하나가 동시에 돌릴 태스크 수입니다. 스레드 하나가 여러 태스크를 맡을 수도 있습니다. 스레드 총합 > 태스크 수 → 남는 스레드는 유휴 상태가 됩니다. 인스턴스를 늘려도 파티션 수를 넘는 병렬성은 얻지 못합니다. 처리량을 더 늘리려면 파티션 수를 늘려야 합니다. 스레드끼리 공유하는 상태가 없으므로 스레드 간 조정도 필요하지 않습니다.
태스크와 스레드 병렬성 — 입력 파티션 수가 태스크 수를 정하고, num.stream.threads와 인스턴스 수가 그 태스크를 어떻게 나눠 갖는지
병렬성·리밸런스 관련 Streams 설정 (Apache Kafka 4.3 기본값)
설정기본값설명튜닝 포인트
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 연산의 상태 저장소에 정확히 한 번 반영됩니다. 입력 토픽 오프셋 커밋, 상태 저장소 갱신, 출력 토픽 쓰기가 원자적으로 완료됩니다.

processing.guarantee와 그것이 바꾸는 것들
항목at_least_once (기본)exactly_once_v2
commit.interval.ms30000 (30초)기본값이 100ms로 바뀝니다
컨슈머 isolation.levelread_uncommittedread_committed로 설정됩니다
프로듀서 enable.idempotence(프로듀서 기본값)true로 설정됩니다
브로커 요구사항2.5.0 이상. 기본 설정으로는 브로커 3대 이상이 필요합니다

인터랙티브 쿼리

Streams는 상태 저장소를 읽기 전용으로 직접 조회하는 것을 허용합니다. 저장소를 만든 애플리케이션 외부의 메서드·스레드·프로세스·애플리케이션에서도 접근할 수 있으며, 모든 저장소는 이름을 가지고 인터랙티브 쿼리는 읽기 연산만 노출합니다.

이 기능의 가치는 출력 토픽을 거치지 않고 최신 집계 결과를 바로 서빙할 수 있다는 점입니다. 집계 결과를 다시 외부 DB에 넣고 그것을 조회하는 단계가 사라집니다. 다만 상태는 인스턴스에 분산되어 있으므로, 찾는 키가 다른 인스턴스에 있으면 그 인스턴스로 요청을 라우팅해야 합니다. 그 메타데이터를 Streams가 제공합니다.

Streams Rebalance Protocol (KIP-1071)

KIP-848이 일반 컨슈머의 리밸런스 조정을 클라이언트에서 브로커로 옮겼듯이, KIP-1071은 같은 모델을 Kafka Streams 워크로드로 확장합니다. 리밸런스 이벤트마다 클라이언트가 배정을 계산하는 대신, 브로커에서 배정이 지속적으로 계산됩니다. 애플리케이션은 컨슈머 그룹이 아니라 streams group으로 등록됩니다.

4.3 기준 streams 프로토콜에서 지원되는 것과 아직 안 되는 것
지원됨아직 지원 안 됨
  • 핵심 streams group 리밸런스 프로토콜 (group.protocol=streams)
  • sticky task assignor (태스크 이동 최소화)
  • 인터랙티브 쿼리
  • StreamsGroupDescribe RPC와 Admin API
  • bin/kafka-streams-groups.sh로 list/describe/delete
  • 오프라인 마이그레이션 (classic ↔ streams)
  • static membershipinstance.id 설정이 거부됩니다
  • 토폴로지 변경 — 소스 토픽 추가나 sub-topology 수 변경 같은 큰 변경 시 새 streams group을 만들어야 합니다
  • High Availability assignor — sticky만 지원되므로 warmup 태스크와 랙 인식 배정이 안 됩니다
  • 정규식 토픽 구독
  • 온라인 마이그레이션 — 실행 중 전환 불가. 정비 시간이 필요합니다
  • 커스텀 KafkaClientSupplier의 "main" 컨슈머 제공

streams group에는 컨슈머 그룹에 없는 NOT_READY 상태가 있습니다. 그룹 코디네이터가 토폴로지에 필요한 소스·내부 토픽이 없거나 설정이 맞지 않는 것을 감지하면 이 상태로 들어가고, 모든 멤버는 빈 배정을 받습니다. 코디네이터가 하는 검사에는 co-partition 그룹이 실제로 co-partition되어 있는지 확인하는 단계가 포함됩니다 — 앞에서 본 요건이 브로커 쪽에서도 검증되기 시작한 것입니다.

테스팅 — 브로커 없이 토폴로지를 검증합니다

CCDAK에는 Application Testing 도메인이 별도로 있고(8%), 그 안에서 가장 자주 다뤄지는 것이 TopologyTestDriver입니다. kafka-streams-test-utils 아티팩트를 테스트 의존성으로 추가하면 씁니다.

pom.xml — 테스트 유틸리티 의존성
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams-test-utils</artifactId>
    <version>4.3.0</version>
    <scope>test</scope>
</dependency>

TopologyTestDriver는 라이브러리 런타임을 시뮬레이션합니다 — 입력 토픽에서 레코드를 계속 가져와 토폴로지를 순회하며 처리하는 동작을 대신하고, 결과 레코드를 캡처하며 내장 상태 저장소를 조회할 수 있게 해 줍니다. 브로커가 필요하지 않습니다.

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로 등록합니다.

어떤 상황에 무엇을 쓰는가

Kafka 테스팅 도구 선택 기준. Apache Kafka 제공외부 도구를 구분했습니다.
목적 도구 브로커 필요 출처
프로듀서 로직 단위 테스트 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에 대응합니다.

ksqlDB — 스트림/테이블 선언과 파생 스트림 생성 (개념 예시)
-- 기존 토픽을 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의 선택 기준
기준Kafka StreamsksqlDB
개발 언어Java / ScalaSQL
배포 형태내 애플리케이션 안에 임베드ksqlDB 서버 클러스터에 질의를 등록
표현력임의 로직 가능 (Processor API 포함)SQL로 표현되는 범위. UDF로 확장
운영 주체애플리케이션 팀ksqlDB 서버를 운영하는 팀
언제 쓰나기존 서비스에 스트림 처리를 넣을 때, 복잡한 로직·외부 연동이 필요할 때탐색적 분석, 빠른 프로토타이핑, SQL에 익숙한 팀이 파이프라인을 직접 만들 때
Apache Kafka 포함포함미포함 (Confluent)

무엇을 모니터링해야 하는가

Kafka Streams 핵심 JMX 메트릭 (Apache Kafka 4.3 모니터링 문서 기준)
메트릭MBean 그룹무엇을 알려 주는가
failed-stream-threadskafka.streams:type=stream-metrics클라이언트 시작 이후 실패한 스트림 스레드 수. 0이 아니면 처리 용량이 줄어든 상태입니다
alive-stream-threadskafka.streams:type=stream-metrics현재 살아 있는(또는 리밸런스 참여 중인) 스레드 수
client-state / thread-statestream-metrics / stream-thread-metrics클라이언트·스레드 상태의 enum ordinal 값
process-rate, process-latency-avgkafka.streams:type=stream-thread-metrics초당 처리 레코드 수와 평균 처리 시간
commit-latency-avgkafka.streams:type=stream-thread-metrics이 스레드의 모든 태스크에 대한 평균 커밋 실행 시간(ms)
poll-latency-avgkafka.streams:type=stream-thread-metrics컨슈머 폴링 평균 시간
blocked-time-ns-totalkafka.streams:type=stream-thread-metrics브로커 때문에 블록된 총 시간. 병목이 내 코드인지 브로커인지 가르는 지표입니다
task-created-ratekafka.streams:type=stream-thread-metrics초당 태스크 생성률. 계속 높으면 리밸런스 스톰입니다
record-e2e-latency-maxstream-processor-node-metrics / stream-state-metrics레코드 타임스탬프와 처리 시점 시스템 시각을 비교한 종단 지연. recording level이 info여서 프로덕션에서도 켜 둘 수 있습니다
restore-ratekafka.streams:type=stream-state-metrics상태 저장소 복원 속도. 재시작 후 얼마나 남았는지 추정에 씁니다

흔한 오해

시험 포인트 정리

확인 문제

네 가지 문항 유형(단일 선택 · 복수 선택 · 연결형 · 순서 배열)이 섞여 있습니다. 리파티션 유발 여부와 co-partitioning 요건은 반드시 맞히고 넘어가세요.

공식 문서 출처

이 장의 설정 기본값·API·의미론은 모두 아래에서 확인했습니다 (Apache Kafka 4.3 문서 기준).