실무 예제 · 9
Kafka Streams 실시간 집계
Kafka Streams는 별도 클러스터가 아니라 애플리케이션에 링크되는 라이브러리입니다. 컨슈머·프로듀서·로컬 상태 저장소를 하나로 묶어 주는 대신, 상태 저장소·changelog·윈도우 시간 개념을 이해하지 않으면 결과가 왜 그렇게 나오는지 설명할 수 없습니다. 이 예제는 1분 단위 클릭 집계를 만들면서 그 부분을 확인합니다.
학습 목표
- tumbling 윈도우 집계를 DSL로 작성하고 grace period가 늦은 이벤트에 어떻게 작용하는지 설명할 수 있습니다.
- 상태 저장소(RocksDB)와 changelog 토픽의 관계, 복구 절차를 알 수 있습니다.
suppress로 중간 결과를 억제해 윈도우당 한 번만 결과를 내보낼 수 있습니다.processing.guarantee=exactly_once_v2가 바꾸는 설정과 트레이드오프를 알 수 있습니다.
시나리오
커머스 사이트의 클릭 스트림을 실시간으로 집계해 상품별 1분간 클릭 수를 대시보드에 표시합니다. 피크 시 초당 8천 건, 상품은 12만 종입니다.
기존에는 5분마다 도는 배치가 전체를 다시 집계했습니다. 문제는 두 가지였습니다 — 지연이 최대 5분이고, 모바일 앱의 오프라인 큐에서 늦게 도착하는 이벤트가 누락되었습니다. 어떤 이벤트는 30초, 어떤 이벤트는 3분 늦게 옵니다.
Streams로 옮기면서 늦게 온 이벤트를 어디까지 받아 줄지를 명시적인 설정(grace period)으로 정합니다.
아키텍처
집계는 상태를 필요로 합니다.
Streams는 그 상태를 로컬 RocksDB에 두고,
같은 내용을 changelog 토픽({application.id}-{store}-changelog)에
compact 정책으로 복제합니다.
인스턴스가 죽으면 다른 인스턴스가 changelog를 재생해 상태를 복원합니다.
num.stream.threads에 따른 배치
사전 요구사항
- 예제 1의 3노드 KRaft 클러스터
- 토픽
clickstream(입력, 파티션 6·RF3),clicks-per-minute(출력, 파티션 6·RF3) — 예제 1의create-topics.sh가 만듭니다 org.apache.kafka:kafka-streams4.3.1, Java 17 이상
전체 코드
clicks-aggregator/
├── pom.xml
└── src/main/java/com/example/streams/
├── ClickCountTopology.java # 토폴로지 정의 (테스트 가능하게 분리)
├── ClickTimestampExtractor.java # 이벤트 시각 추출
└── Main.java # 설정 + 기동 + 종료
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>clicks-aggregator</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.release>17</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<kafka.version>4.3.1</kafka.version>
<slf4j.version>2.0.17</slf4j.version>
</properties>
<dependencies>
<!-- kafka-streams 는 kafka-clients 를 전이 의존성으로 가져옵니다. -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>${kafka.version}</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>${slf4j.version}</version>
</dependency>
<!-- TopologyTestDriver: 브로커 없이 토폴로지를 단위 테스트합니다. -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams-test-utils</artifactId>
<version>${kafka.version}</version>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>3.5.0</version>
<configuration>
<mainClass>com.example.streams.Main</mainClass>
</configuration>
</plugin>
</plugins>
</build>
</project>
이벤트 시각 추출
package com.example.streams;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.streams.processor.TimestampExtractor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
/**
* 이벤트 페이로드의 clickedAt 필드를 이벤트 시각으로 씁니다.
*
* 기본 추출기는 레코드의 Kafka 타임스탬프를 쓰는데,
* 그것은 "프로듀서가 발행한 시각"(CreateTime) 또는
* "브로커가 기록한 시각"(LogAppendTime)입니다.
* 모바일 오프라인 큐처럼 발행이 늦어지는 경우
* 실제 클릭 시각과 크게 벌어지므로 페이로드에서 직접 꺼냅니다.
*/
public class ClickTimestampExtractor implements TimestampExtractor {
private static final Logger log = LoggerFactory.getLogger(ClickTimestampExtractor.class);
private static final Pattern CLICKED_AT = Pattern.compile("\"clickedAt\"\\s*:\\s*(\\d+)");
@Override
public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
Object value = record.value();
if (value instanceof String json) {
Matcher m = CLICKED_AT.matcher(json);
if (m.find()) {
return Long.parseLong(m.group(1));
}
}
// 추출 실패 시 음수를 반환하면 Streams 가 그 레코드를 버립니다.
// 여기서는 파티션의 현재 시각(스트림 시각)으로 대체해 유실을 막습니다.
// 어느 쪽이 맞는지는 도메인 판단입니다 — 조용히 버리면 원인을 못 찾습니다.
log.warn("clickedAt 추출 실패 — partitionTime 으로 대체합니다. topic={} offset={}",
record.topic(), record.offset());
return partitionTime >= 0 ? partitionTime : record.timestamp();
}
}
토폴로지
package com.example.streams;
import org.apache.kafka.common.serialization.Serde;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.kstream.Consumed;
import org.apache.kafka.streams.kstream.Grouped;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.kstream.Suppressed;
import org.apache.kafka.streams.kstream.TimeWindows;
import org.apache.kafka.streams.kstream.Windowed;
import org.apache.kafka.streams.state.Stores;
import org.apache.kafka.streams.state.WindowBytesStoreSupplier;
import java.time.Duration;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
/**
* 상품별 1분 클릭 수를 집계합니다.
*
* 토폴로지를 Main 에서 분리해 두면 TopologyTestDriver 로
* 브로커 없이 단위 테스트할 수 있습니다.
*/
public final class ClickCountTopology {
public static final String STORE_NAME = "clicks-per-minute-store";
private static final Pattern PRODUCT_ID = Pattern.compile("\"productId\"\\s*:\\s*\"([^\"]+)\"");
/** 윈도우 크기 — 1분 tumbling */
private static final Duration WINDOW_SIZE = Duration.ofMinutes(1);
/**
* grace period — 윈도우가 끝난 뒤에도 이만큼은 늦은 이벤트를 받아 줍니다.
* 시나리오의 "최대 3분 늦게 오는 모바일 이벤트" 를 근거로 3분을 줍니다.
* 이 값을 늘리면 정확도가 오르지만 결과 확정이 그만큼 늦어지고
* 상태 저장소가 커집니다. 정확도와 지연의 트레이드오프입니다.
*/
private static final Duration GRACE = Duration.ofMinutes(3);
private ClickCountTopology() {
}
public static Topology build(String inputTopic, String outputTopic) {
StreamsBuilder builder = new StreamsBuilder();
Serde<String> stringSerde = Serdes.String();
Serde<Long> longSerde = Serdes.Long();
KStream<String, String> clicks = builder.stream(
inputTopic,
Consumed.with(stringSerde, stringSerde)
// 이벤트 시각을 페이로드에서 꺼냅니다.
.withTimestampExtractor(new ClickTimestampExtractor()));
// 상태 저장소를 명시적으로 만듭니다.
// 윈도우 저장소의 보관 기간은 최소한 (윈도우 크기 + grace) 이상이어야 합니다.
// 짧으면 grace 안에 도착한 이벤트를 받을 저장소가 이미 없어집니다.
WindowBytesStoreSupplier storeSupplier = Stores.persistentWindowStore(
STORE_NAME,
WINDOW_SIZE.plus(GRACE).plus(Duration.ofMinutes(1)), // retention
WINDOW_SIZE, // windowSize
false); // retainDuplicates
clicks
// 키를 productId 로 다시 잡습니다.
// selectKey 는 "리파티션을 유발" 합니다 — 키가 바뀌면 파티션이 바뀌므로
// Streams 가 내부 리파티션 토픽을 자동으로 만듭니다.
// (mapValues 처럼 값만 바꾸는 연산은 리파티션을 유발하지 않습니다)
.selectKey((key, value) -> extractProductId(value))
// 키가 null 인 레코드는 집계할 수 없으므로 걸러 냅니다.
// 걸러 내지 않으면 NullPointerException 으로 스트림 스레드가 죽습니다.
.filter((productId, value) -> productId != null)
.groupByKey(Grouped.with("clicks-by-product", stringSerde, stringSerde))
// ofSizeAndGrace 가 4.x 의 표준 API 입니다.
// (구버전의 TimeWindows.of(...) + until(...) 조합은 제거되었습니다)
.windowedBy(TimeWindows.ofSizeAndGrace(WINDOW_SIZE, GRACE))
.count(Materialized.<String, Long>as(storeSupplier)
.withKeySerde(stringSerde)
.withValueSerde(longSerde))
// suppress: 윈도우가 닫힐 때까지 중간 결과를 내보내지 않습니다.
// 이것이 없으면 이벤트가 들어올 때마다 갱신된 카운트가 계속 발행되어
// 하류 대시보드가 1, 2, 3, 4 ... 를 모두 보게 됩니다.
// untilWindowCloses 는 StrictBufferConfig 만 받습니다 —
// 버퍼가 넘칠 때 결과를 조기 방출하면 "윈도우당 한 번" 보장이 깨지기 때문입니다.
.suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
.toStream()
// Windowed<String> 키를 "productId@윈도우시작시각" 문자열로 펼칩니다.
// 하류가 일반 컨슈머라면 윈도우 Serde 를 알 필요가 없어집니다.
.map((Windowed<String> windowedKey, Long count) -> KeyValue.pair(
windowedKey.key(),
"{\"productId\":\"%s\",\"windowStart\":%d,\"windowEnd\":%d,\"clicks\":%d}"
.formatted(windowedKey.key(),
windowedKey.window().start(),
windowedKey.window().end(),
count)))
.to(outputTopic, Produced.with(stringSerde, stringSerde));
return builder.build();
}
private static String extractProductId(String json) {
if (json == null) {
return null;
}
Matcher m = PRODUCT_ID.matcher(json);
return m.find() ? m.group(1) : null;
}
}
설정과 기동
package com.example.streams;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;
public final class Main {
private static final Logger log = LoggerFactory.getLogger(Main.class);
private static final String BOOTSTRAP =
"localhost:29092,localhost:39092,localhost:49092";
private static final String INPUT = "clickstream";
private static final String OUTPUT = "clicks-per-minute";
public static void main(String[] args) {
boolean eos = args.length > 0 && "eos".equals(args[0]);
Properties p = new Properties();
// application.id 는 이 Streams 애플리케이션의 정체성입니다.
// - 컨슈머 그룹 ID 가 됩니다.
// - 내부 토픽(changelog, repartition) 이름의 접두어가 됩니다.
// - 로컬 상태 디렉터리 경로에 들어갑니다.
// 바꾸면 완전히 새 애플리케이션이 되어 상태를 처음부터 다시 만듭니다.
p.put(StreamsConfig.APPLICATION_ID_CONFIG, "clicks-aggregator");
p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 한 인스턴스가 돌릴 스트림 스레드 수. 기본값은 1 입니다.
// 태스크 수(=입력 파티션 수 6)를 넘기면 남는 스레드는 유휴 상태가 됩니다.
p.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 3);
// 내부 토픽(changelog, repartition)의 복제 계수.
// 기본값은 -1 로 "브로커 기본값을 따름" 을 의미하며,
// 브로커 default.replication.factor 가 1 이면 상태가 단일 사본이 됩니다.
// 반드시 명시하세요.
p.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);
// 로컬 상태 디렉터리. 컨테이너에서는 반드시 영속 볼륨을 마운트해야 합니다.
// 휘발성이면 재시작마다 changelog 전체를 재생해 복구가 오래 걸립니다.
p.put(StreamsConfig.STATE_DIR_CONFIG, "/tmp/kafka-streams");
// standby replica: 다른 인스턴스가 상태 사본을 미리 따라 만들어 둡니다.
// 기본값 0. 1 이상으로 두면 인스턴스 장애 시 복구가 훨씬 빠릅니다.
// 대가는 changelog 재생 트래픽과 디스크입니다.
p.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);
if (eos) {
// processing.guarantee 의 기본값은 at_least_once 이고,
// 유효값은 at_least_once / exactly_once_v2 두 개뿐입니다
// (exactly_once 와 exactly_once_beta 는 4.0 에서 제거되었습니다).
//
// 켜면 commit.interval.ms 기본값이 30000 -> 100 으로 바뀝니다.
// 커밋이 잦아져 처리량이 떨어지고 브로커 요청이 늘어나는 대신,
// 입력 오프셋과 출력·상태 갱신이 원자적으로 커밋됩니다.
p.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
log.info("processing.guarantee = exactly_once_v2");
} else {
p.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.AT_LEAST_ONCE);
// at_least_once 의 기본 커밋 주기는 30000 입니다.
// 대시보드 지연을 줄이려고 짧게 둡니다.
p.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 5_000);
}
Topology topology = ClickCountTopology.build(INPUT, OUTPUT);
// 토폴로지를 로그로 남기면 리파티션 토픽이 어디 생기는지 눈으로 확인할 수 있습니다.
log.info("토폴로지:\n{}", topology.describe());
KafkaStreams streams = new KafkaStreams(topology, p);
// 스트림 스레드에서 잡히지 않은 예외가 났을 때의 동작을 정합니다.
// 기본 동작은 스레드 종료이며, 모든 스레드가 죽으면 애플리케이션이 멈춥니다.
streams.setUncaughtExceptionHandler((Throwable e) -> {
log.error("스트림 스레드에서 처리되지 않은 예외", e);
// REPLACE_THREAD : 그 스레드만 교체하고 계속 처리합니다.
// SHUTDOWN_CLIENT : 이 인스턴스만 종료합니다.
// SHUTDOWN_APPLICATION : 같은 application.id 의 모든 인스턴스를 종료합니다.
//
// 데이터 오류로 인한 예외를 REPLACE_THREAD 로 두면
// 같은 레코드에서 무한히 실패하는 루프가 됩니다.
// 원인을 구분할 수 없다면 SHUTDOWN_CLIENT 가 안전합니다.
return StreamsUncaughtExceptionHandler
.StreamThreadExceptionResponse.SHUTDOWN_CLIENT;
});
streams.setStateListener((newState, oldState) ->
log.info("상태 전이 {} → {}", oldState, newState));
CountDownLatch latch = new CountDownLatch(1);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("SIGTERM 수신 — Streams 를 정상 종료합니다");
// close(Duration) 은 진행 중인 처리를 마치고 오프셋을 커밋한 뒤 닫습니다.
// 타임아웃 없이 kill 되면 changelog 재생이 필요한 상태로 남습니다.
streams.close(Duration.ofSeconds(60));
latch.countDown();
}, "streams-shutdown-hook"));
try {
streams.start();
latch.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} catch (RuntimeException e) {
log.error("Streams 기동 실패", e);
System.exit(1);
}
log.info("정상 종료");
}
}
실행 방법
# 0. 예제 1의 클러스터와 토픽
cd kafka-lab && docker compose ps
./kcli kafka-topics.sh --describe --topic clickstream | head -2
# 1. 빌드 후 기동 (at_least_once)
cd ../clicks-aggregator
mvn -q clean package
mvn -q exec:java
# 2. 다른 터미널에서 클릭 이벤트를 넣습니다.
# clickedAt 을 "지금" 으로 넣어야 윈도우가 정상 동작합니다.
cd kafka-lab
NOW=$(date +%s000)
for i in $(seq 1 30); do
PID="P-$((i % 3 + 1))"
echo "$PID:{\"productId\":\"$PID\",\"userId\":\"U-$i\",\"clickedAt\":$NOW}"
done | ./kcli kafka-console-producer.sh --topic clickstream \
--property parse.key=true --property key.separator=:
# 3. 결과를 봅니다. 윈도우가 닫히고 grace(3분)가 지난 뒤에 나옵니다.
./kcli kafka-console-consumer.sh --topic clicks-per-minute --from-beginning \
--property print.key=true
검증 방법
1. 내부 토픽이 만들어졌는가
./kcli kafka-topics.sh --list | grep clicks-aggregator
clicks-aggregator-clicks-by-product-repartition
clicks-aggregator-clicks-per-minute-store-changelog
repartition 토픽은 selectKey로 키를 바꿨기 때문에 생겼습니다.
mapValues처럼 값만 바꾸는 연산으로는 생기지 않습니다.
changelog 토픽은 상태 저장소의 백업이며 compact 정책입니다.
compact여야 합니다./kcli kafka-topics.sh --describe \
--topic clicks-aggregator-clicks-per-minute-store-changelog
Topic: clicks-aggregator-clicks-per-minute-store-changelog PartitionCount: 6 ReplicationFactor: 3
Configs: cleanup.policy=compact,delete,min.insync.replicas=2,retention.ms=...
윈도우 저장소의 changelog는 compact,delete 조합을 씁니다 —
키별 최신값을 유지하면서 보관 기간이 지난 윈도우는 삭제합니다.
ReplicationFactor: 3이 나오는 것은
StreamsConfig.REPLICATION_FACTOR_CONFIG=3을 명시했기 때문입니다.
명시하지 않으면 브로커 default.replication.factor(기본 1)를 따릅니다.
2. 집계 결과가 맞는가
P-1 {"productId":"P-1","windowStart":1785024000000,"windowEnd":1785024060000,"clicks":10}
P-2 {"productId":"P-2","windowStart":1785024000000,"windowEnd":1785024060000,"clicks":10}
P-3 {"productId":"P-3","windowStart":1785024000000,"windowEnd":1785024060000,"clicks":10}
30건을 3개 상품에 나눠 넣었으므로 각 10이 나옵니다.
windowStart가 분 경계(00초)에 정렬되어 있는지 확인하세요 —
tumbling 윈도우는 epoch 기준으로 정렬되며,
이벤트가 들어온 시각이 아니라 고정된 경계를 씁니다.
3. grace period가 늦은 이벤트를 받아 주는가
LATE=$(( $(date +%s) * 1000 - 120000 ))
echo "P-1:{\"productId\":\"P-1\",\"userId\":\"U-LATE\",\"clickedAt\":$LATE}" \
| ./kcli kafka-console-producer.sh --topic clickstream \
--property parse.key=true --property key.separator=:
# grace 를 넘는 5분 전 이벤트 — 버려집니다.
TOO_LATE=$(( $(date +%s) * 1000 - 300000 ))
echo "P-1:{\"productId\":\"P-1\",\"userId\":\"U-TOOLATE\",\"clickedAt\":$TOO_LATE}" \
| ./kcli kafka-console-producer.sh --topic clickstream \
--property parse.key=true --property key.separator=:
kafka.streams:type=stream-task-metrics,thread-id=*,task-id=*
Attribute: dropped-records-total
Attribute: dropped-records-rate
dropped-records-total이 증가하면 grace period가 짧다는 신호입니다.
이 메트릭에 알림을 걸어 두지 않으면
"집계가 조금씩 부족한데 이유를 모르겠다"는 상태가 됩니다.
늦은 이벤트는 에러 없이 조용히 버려집니다.
4. 상태 복구가 동작하는가
# 애플리케이션을 정상 종료한 뒤 로컬 상태를 지웁니다.
pkill -f 'com.example.streams.Main'
rm -rf /tmp/kafka-streams/clicks-aggregator
# 재시작하면 changelog 를 재생해 상태를 복원합니다.
cd ../clicks-aggregator && mvn -q exec:java
INFO com.example.streams.Main - 상태 전이 CREATED → REBALANCING
INFO o.a.k.s.p.internals.StoreChangelogReader - stream-thread [...] Restoration in progress for 6 partitions.
INFO o.a.k.s.p.internals.StoreChangelogReader - stream-thread [...] Finished restoring changelog clicks-aggregator-clicks-per-minute-store-changelog-0
INFO com.example.streams.Main - 상태 전이 REBALANCING → RUNNING
REBALANCING에 머무는 시간이 복구에 걸리는 시간입니다.
상태가 크면 이 시간이 수십 분까지 늘어납니다 —
num.standby.replicas를 1 이상으로 두면
다른 인스턴스가 미리 상태를 따라 만들어 두므로 이 시간이 크게 줄어듭니다.
5. exactly_once_v2가 커밋 주기를 바꾸는가
mvn -q exec:java -Dexec.args="eos"
commit.interval.ms가 100으로 바뀝니다INFO o.a.k.streams.StreamsConfig - StreamsConfig values:
commit.interval.ms = 100
processing.guarantee = exactly_once_v2
...
INFO o.a.k.c.p.internals.TransactionManager - [Producer clientId=clicks-aggregator-...] ProducerId set to 7000 with epoch 0
ProducerId set to ... 로그가 나오면 트랜잭션 프로듀서가 쓰이는 것입니다.
transactional.id는 Streams가 application.id와 태스크에서 자동으로 만들므로
직접 정하지 않습니다.
하류에서 이 결과를 읽는 컨슈머는 isolation.level=read_committed여야
EOS가 완성됩니다(예제 5).
이 예제에서 쓴 Streams 설정
| 설정 | 기본값 | 이 예제 값 | 이유 |
|---|---|---|---|
application.id | 필수 | clicks-aggregator | 컨슈머 그룹 ID, 내부 토픽 접두어, 상태 디렉터리 이름이 됩니다. 바꾸면 새 애플리케이션입니다 |
num.stream.threads | 1 | 3 | 기본값이 1이라는 점에 주의. 파티션 수(6)가 상한입니다 |
processing.guarantee | at_least_once | 인자로 선택 | 유효값은 at_least_once / exactly_once_v2 두 개뿐입니다 |
commit.interval.ms | 30000(EOS일 때 100) | 5000 (ALO) | EOS를 켜면 기본값이 자동으로 100으로 바뀝니다 |
replication.factor | -1 (브로커 기본값) | 3 | 명시하지 않으면 브로커 default.replication.factor(기본 1)를 따라 상태가 단일 사본이 됩니다 |
num.standby.replicas | 0 | 1 | 장애 복구 시간을 줄입니다. 대가는 changelog 재생 트래픽과 디스크 |
state.dir | ${java.io.tmpdir}/kafka-streams리눅스에서는 보통 /tmp/kafka-streams | 동일 | /tmp는 재부팅 시 사라집니다. 프로덕션에서는 영속 경로로 바꾸세요 |
default.timestamp.extractor | 레코드 타임스탬프 사용 | 커스텀 | 발행 지연이 큰 소스는 페이로드의 이벤트 시각을 써야 윈도우가 맞습니다 |
selectKey·map은 리파티션을 유발하고 mapValues는 유발하지 않습니다
프로덕션 고려사항
| 항목 | 이 예제 | 프로덕션 |
|---|---|---|
| 상태 디렉터리 | /tmp/kafka-streams |
영속 볼륨(K8s라면 StatefulSet + PVC). 휘발성이면 배포마다 changelog 전체를 재생해 복구가 수십 분 걸립니다 |
| 인스턴스 식별 | 없음 | group.instance.id(static membership)로 재시작 시 리밸런스를 회피합니다. 상태가 큰 애플리케이션에서 효과가 가장 큽니다 |
| 직렬화 | Serdes.String() + 정규식 파싱 |
Avro/Protobuf Serde. 정규식 파싱은 스키마 변경에 조용히 깨집니다(예제 7) |
| 버려진 레코드 | 메트릭만 존재 | dropped-records-total에 알림을 걸어야 합니다. grace를 넘긴 이벤트는 에러 없이 사라집니다 |
| 예외 처리 | SHUTDOWN_CLIENT |
역직렬화 오류는 default.deserialization.exception.handler로, 처리 오류는 default.processing.exception.handler로 분리해 데이터 오류와 인프라 오류를 다르게 다룹니다 |
| RocksDB 튜닝 | 기본값 | 상태가 크면 블록 캐시·write buffer를 RocksDBConfigSetter로 조정합니다. RocksDB는 힙 밖 메모리를 쓰므로 컨테이너 메모리 한계 계산에 반드시 포함해야 합니다 |
| 내부 토픽 | 자동 생성 | 자동 생성되지만 RF는 replication.factor 설정을 따릅니다. 명시하지 않으면 브로커 기본값(1)이 되어 상태가 단일 사본이 됩니다 |
| EOS | 인자로 선택 | 켜면 처리량이 떨어집니다. 하류 컨슈머가 read_committed가 아니면 의미가 없습니다 |
| 토폴로지 변경 | 자유롭게 | 토폴로지를 바꾸면 내부 토픽 이름과 태스크 배정이 달라져 호환되지 않는 경우가 있습니다. 그때는 application.id를 바꾸고 상태를 새로 만드는 것이 안전합니다 |
자주 하는 실수
관련 케이스 스터디
이어서 볼 곳
공식 문서 출처
- Kafka Streams Configs —
num.stream.threads=1,processing.guarantee=at_least_once(유효값at_least_once/exactly_once_v2),commit.interval.ms=30000이며 EOS일 때100,replication.factor=-1,num.standby.replicas=0,state.dir - Streams DSL —
TimeWindows.ofSizeAndGrace/ofSizeWithNoGrace,Suppressed.untilWindowCloses(StrictBufferConfig만 허용),Suppressed.untilTimeLimit,BufferConfig.unbounded()/maxRecords/maxBytes, 리파티션 유발 연산 목록 - Streams Architecture — 태스크 수 = 입력 파티션 수, 스레드에 태스크가 배분되는 모델, 상태 저장소와 changelog
- Managing Streams Application Topics — 내부 토픽 명명 규칙(
{application.id}-{store}-changelog,-repartition)과 정책 - Streams Task Metrics —
dropped-records-total/dropped-records-rate - StreamThreadExceptionResponse —
REPLACE_THREAD/SHUTDOWN_CLIENT/SHUTDOWN_APPLICATION - Broker Configs —
default.replication.factor=1