학습 목표

시나리오

커머스 사이트의 클릭 스트림을 실시간으로 집계해 상품별 1분간 클릭 수를 대시보드에 표시합니다. 피크 시 초당 8천 건, 상품은 12만 종입니다.

기존에는 5분마다 도는 배치가 전체를 다시 집계했습니다. 문제는 두 가지였습니다 — 지연이 최대 5분이고, 모바일 앱의 오프라인 큐에서 늦게 도착하는 이벤트가 누락되었습니다. 어떤 이벤트는 30초, 어떤 이벤트는 3분 늦게 옵니다.

Streams로 옮기면서 늦게 온 이벤트를 어디까지 받아 줄지를 명시적인 설정(grace period)으로 정합니다.

아키텍처

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 → 리파티션 → 윈도우 집계 → suppress → sink와 sub-topology 경계
상태 저장소와 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로 복구하는 과정

집계는 상태를 필요로 합니다. Streams는 그 상태를 로컬 RocksDB에 두고, 같은 내용을 changelog 토픽({application.id}-{store}-changelog)에 compact 정책으로 복제합니다. 인스턴스가 죽으면 다른 인스턴스가 changelog를 재생해 상태를 복원합니다.

태스크와 스레드 병렬성 — 파티션 수가 병렬성의 상한입니다 입력 토픽의 파티션 여섯 개가 태스크 여섯 개로 이어지고, 그 태스크가 애플리케이션 인스턴스 두 대의 스트림 스레드 세 개씩에 하나씩 배정된 모습을 그린 그림입니다. 태스크 수는 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에 따른 배치
윈도우 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 윈도우를 적용한 결과

사전 요구사항

전체 코드

디렉터리 구조
clicks-aggregator/
├── pom.xml
└── src/main/java/com/example/streams/
    ├── ClickCountTopology.java   # 토폴로지 정의 (테스트 가능하게 분리)
    ├── ClickTimestampExtractor.java # 이벤트 시각 추출
    └── Main.java                 # 설정 + 기동 + 종료
clicks-aggregator/pom.xml
<?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>

이벤트 시각 추출

src/main/java/com/example/streams/ClickTimestampExtractor.java
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();
    }
}

토폴로지

src/main/java/com/example/streams/ClickCountTopology.java
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;
    }
}

설정과 기동

src/main/java/com/example/streams/Main.java
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. 내부 토픽이 만들어졌는가

changelog와 repartition 토픽 확인
./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 정책입니다.

changelog 정책 확인 — 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가 늦은 이벤트를 받아 주는가

2분 전 시각의 이벤트를 넣습니다 (grace 3분 안)
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가 커밋 주기를 바꾸는가

EOS 모드로 실행
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 설정

Apache Kafka 4.3 StreamsConfig
설정기본값이 예제 값이유
application.id필수clicks-aggregator컨슈머 그룹 ID, 내부 토픽 접두어, 상태 디렉터리 이름이 됩니다. 바꾸면 새 애플리케이션입니다
num.stream.threads13기본값이 1이라는 점에 주의. 파티션 수(6)가 상한입니다
processing.guaranteeat_least_once인자로 선택유효값은 at_least_once / exactly_once_v2 두 개뿐입니다
commit.interval.ms30000
(EOS일 때 100)
5000 (ALO)EOS를 켜면 기본값이 자동으로 100으로 바뀝니다
replication.factor-1 (브로커 기본값)3명시하지 않으면 브로커 default.replication.factor(기본 1)를 따라 상태가 단일 사본이 됩니다
num.standby.replicas01장애 복구 시간을 줄입니다. 대가는 changelog 재생 트래픽과 디스크
state.dir${java.io.tmpdir}/kafka-streams
리눅스에서는 보통 /tmp/kafka-streams
동일/tmp는 재부팅 시 사라집니다. 프로덕션에서는 영속 경로로 바꾸세요
default.timestamp.extractor레코드 타임스탬프 사용커스텀발행 지연이 큰 소스는 페이로드의 이벤트 시각을 써야 윈도우가 맞습니다
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 / 리파티션 유발 연산 구분 — 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를 바꾸고 상태를 새로 만드는 것이 안전합니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처