Step 13 — Kafka Streams

학습 목표

  • Streams 가 별도 클러스터가 아니라 Consumer + Producer 위에 얹힌 라이브러리임을 설명한다
  • topology.describe() 출력을 읽고 서브토폴로지 경계를 판별한다
  • KStream / KTable / GlobalKTable 의 차이와 조인 가능 조합을 표로 정리한다
  • selectKey 뒤에 repartition 토픽이, 집계 뒤에 changelog 토픽이 자동 생성되는 것을 kt --list 로 직접 확인한다
  • 윈도우 4종을 비교하고, grace period 를 지나 도착한 레코드가 조용히 버려지는 것을 재현한다
  • co-partitioning 위반으로 TopologyException 이 나는 것을 재현한다
  • kafka-streams-application-reset.sh 로 내부 토픽과 상태 저장소를 정리한다

선행 스텝: Step 12 — Kafka Connect 예상 소요: 120분


13-0. 실습 준비

이 스텝은 Practice.java 하나로 진행합니다. 시나리오를 인자로 골라 실행합니다.

cd docs/reference/kafka/step-13-streams
docker cp Practice.java kafka-1:/tmp/
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java'

결과

usage: java -cp "/opt/kafka/libs/*" Practice.java <scenario>

  topology   토폴로지 설명만 출력하고 종료 (실행하지 않음)
  stateless  filter / map / mapValues / flatMap / branch / merge / peek
  selectkey  selectKey → repartition 토픽이 생기는 것을 관찰
  count      groupByKey().count() → changelog 토픽과 상태 저장소
  window     1분 텀블링 윈도우 집계 + grace period
  suppress   suppress(untilWindowCloses) 로 최종 결과만
  join       orders × payments 윈도우 조인
  copartition co-partitioning 위반 재현 (TopologyException)
  eos        processing.guarantee=exactly_once_v2

  예) java -cp "/opt/kafka/libs/*" Practice.java count

/opt/kafka/libs/kafka-streams-3.7.1.jar 이 이미 들어 있으므로 별도 빌드 도구가 필요 없습니다. 확인해 봅니다.

docker exec kafka-1 ls /opt/kafka/libs/ | grep -E 'kafka-streams|rocksdb'

결과

kafka-streams-3.7.1.jar
kafka-streams-examples-3.7.1.jar
rocksdbjni-7.9.2.jar

rocksdbjni 가 상태 저장소의 실체입니다. 13-6 에서 다시 봅니다.

실습 토픽을 만듭니다.

kt --create --topic s13_orders   --partitions 3 --replication-factor 3
kt --create --topic s13_payments --partitions 3 --replication-factor 3
kt --create --topic s13_out      --partitions 3 --replication-factor 3
kt --create --topic s13_left     --partitions 3 --replication-factor 3
kt --create --topic s13_right    --partitions 6 --replication-factor 3

결과

Created topic s13_orders.
Created topic s13_payments.
Created topic s13_out.
Created topic s13_left.
Created topic s13_right.

s13_right 만 파티션이 6개입니다. 13-9 의 co-partitioning 함정을 위한 의도적 설정입니다.

시작 시점의 토픽 목록을 기록해 둡니다. 13-7 에서 before/after 를 비교합니다.

kt --list > /tmp/topics-before.txt
cat /tmp/topics-before.txt

결과

__consumer_offsets
dlq
order-events
orders
payments
s13_left
s13_orders
s13_out
s13_payments
s13_right

13-1. Streams 는 라이브러리입니다

Kafka Streams 를 처음 보면 "Spark 나 Flink 같은 처리 클러스터"를 떠올리기 쉽습니다. 아닙니다.

Spark / FlinkKafka Streams
배포 형태클러스터 + 잡 제출일반 JAR. java -jar 로 실행
리소스 관리YARN / K8s / 자체 매니저없음. 프로세스가 전부
확장클러스터에 노드 추가같은 application.id 로 프로세스를 하나 더 띄움
상태 저장체크포인트 (HDFS/S3)로컬 RocksDB + changelog 토픽
장애 복구잡 매니저가 재실행컨슈머 그룹 리밸런싱
의존성별도 클러스터Kafka 뿐

Streams 애플리케이션의 실체를 벗겨 보면 이렇습니다.

   여러분의 Streams 앱 (평범한 JVM 프로세스)
   ┌───────────────────────────────────────────────────────────┐
   │  StreamThread 1                StreamThread 2             │
   │  ┌──────────────────┐          ┌──────────────────┐       │
   │  │ Consumer         │          │ Consumer         │       │
   │  │   ↓ poll()       │          │   ↓ poll()       │       │
   │  │ Task 0_0         │          │ Task 0_1         │       │
   │  │   → 처리 로직     │          │   → 처리 로직     │       │
   │  │   → RocksDB      │          │   → RocksDB      │       │
   │  │   ↓              │          │   ↓              │       │
   │  │ Producer         │          │ Producer         │       │
   │  └──────────────────┘          └──────────────────┘       │
   └───────────────────────────────────────────────────────────┘
             │                              │
             ▼                              ▼
      Kafka: 입력 토픽 / changelog / repartition / 출력 토픽

컨슈머 그룹 ID 는 application.id 와 같습니다. 그래서 Step 05 에서 배운 것이 전부 그대로 적용됩니다.

kcg --list

결과 (앱을 띄운 뒤)

s13-count-app

kcg --describe --group s13-count-app 로 랙도 그대로 확인됩니다. Streams 앱의 랙 모니터링은 일반 컨슈머와 완전히 동일합니다.

💡 실무 팁 — 확장은 프로세스를 더 띄우는 것뿐입니다 파티션이 3개인 토픽을 처리하는 Streams 앱은 태스크가 3개 생깁니다. 프로세스 1개면 3개를 다 처리하고, 프로세스를 3개 띄우면 각각 1개씩 나눠 갖습니다. 코드도 설정도 안 바꿉니다. 프로세스 4개를 띄우면 하나는 놉니다. Step 05 의 "컨슈머 수 > 파티션 수이면 논다"가 그대로입니다. num.stream.threads 로 한 프로세스 안의 스레드 수를 늘릴 수도 있습니다. 프로세스 1개 × 스레드 3개와 프로세스 3개 × 스레드 1개는 처리 병렬성 면에서 같습니다.


13-2. Topology — 무엇을 실행할지의 설계도

Streams DSL 로 쓴 코드는 실행 전에 Topology(처리 그래프) 로 컴파일됩니다. describe() 로 볼 수 있습니다.

StreamsBuilder b = new StreamsBuilder();
b.stream("orders", Consumed.with(Serdes.String(), Serdes.String()))
 .filter((k, v) -> v.contains("CREATED"))
 .mapValues(v -> v.toUpperCase())
 .to("s13_out", Produced.with(Serdes.String(), Serdes.String()));

Topology topology = b.build();
System.out.println(topology.describe());
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java topology'

결과

Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [orders])
      --> KSTREAM-FILTER-0000000001
    Processor: KSTREAM-FILTER-0000000001 (stores: [])
      --> KSTREAM-MAPVALUES-0000000002
      <-- KSTREAM-SOURCE-0000000000
    Processor: KSTREAM-MAPVALUES-0000000002 (stores: [])
      --> KSTREAM-SINK-0000000003
      <-- KSTREAM-FILTER-0000000001
    Sink: KSTREAM-SINK-0000000003 (topic: s13_out)
      <-- KSTREAM-MAPVALUES-0000000002

읽는 법:

표기
Sub-topology: N독립적으로 실행되는 단위. 각각 별도 태스크가 됨
Source:토픽에서 읽는 노드
Processor:변환 노드. (stores: [...]) 에 쓰는 상태 저장소가 나옴
Sink:토픽에 쓰는 노드
-->다음 노드
<--이전 노드
0000000000노드 순번. DSL 호출 순서대로 매겨짐

서브토폴로지가 하나뿐이면 데이터가 네트워크를 거치지 않습니다. 소스에서 싱크까지 한 스레드 안에서 처리됩니다. 서브토폴로지가 나뉘는 순간(13-5) 그 사이에 repartition 토픽이 끼어들고, 데이터가 Kafka 를 한 번 왕복합니다.

💡 실무 팁 — describe() 를 로그에 남기세요 운영 중인 Streams 앱이 왜 느린지 조사할 때 가장 먼저 보는 것이 토폴로지입니다. 서브토폴로지가 몇 개인지가 곧 네트워크 왕복 횟수이기 때문입니다. 앱 기동 시 log.info("{}", topology.describe()) 를 한 줄 넣어 두면, 코드 변경이 토폴로지를 어떻게 바꿨는지 배포마다 확인할 수 있습니다. https://zz85.github.io/kafka-streams-viz/ 에 이 출력을 붙여 넣으면 그림으로도 볼 수 있습니다.


13-3. KStream vs KTable vs GlobalKTable

세 추상화의 차이가 Streams 학습의 절반입니다.

같은 입력을 셋으로 해석해 봅니다. 입력은 이렇습니다.

key=C001, value=100
key=C002, value=200
key=C001, value=300
KStreamKTableGlobalKTable
해석레코드 스트림 — 독립된 사건 3개변경로그 스트림 — 현재 상태KTable 과 같으나 복제 방식이 다름
결과C001=100, C002=200, C001=300 (3건)C001=300, C002=200 (2건)좌동
비유은행 입출금 내역은행 잔액모든 지점이 갖고 있는 잔액 사본
null 값그냥 값이 null 인 레코드삭제 (tombstone)삭제
파티셔닝파티션별로 나뉨파티션별로 나뉨모든 인스턴스가 전체를 복제
크기 제한없음파티션 크기인스턴스 메모리/디스크에 다 들어가야 함
만드는 법builder.stream(...)builder.table(...)builder.globalTable(...)

조인 가능 조합

왼쪽오른쪽가능?조건
KStreamKStream윈도우 필수. co-partition 필요
KStreamKTable윈도우 없음. co-partition 필요
KStreamGlobalKTableco-partition 불필요. 키가 달라도 됨
KTableKTableco-partition 필요
KTableKStream방향을 뒤집어 KStream-KTable 로
KTableGlobalKTable지원 안 함
GlobalKTable무엇이든GlobalKTable 은 조인의 오른쪽에만

💡 GlobalKTable 은 "작고 잘 안 변하는 참조 데이터"용입니다 상품 카테고리 코드표, 국가 코드, 환율 같은 것입니다. 모든 인스턴스가 전체를 복제하므로 co-partitioning 이 필요 없고, 키가 달라도 KeyValueMapper 로 매핑해서 조인할 수 있습니다. 대신 인스턴스가 10대면 같은 데이터를 10벌 갖습니다. 주문 테이블 같은 것을 GlobalKTable 로 만들면 앱이 뜨지도 못합니다. "이게 각 인스턴스 디스크에 다 들어가는가"가 유일한 판단 기준입니다.


13-4. 스테이트리스 연산

상태를 안 쓰는 연산들입니다. 레코드 하나만 보고 처리하므로 빠르고, 내부 토픽도 안 만듭니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java stateless' &
sleep 15
docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_orders \
  --property parse.key=true --property key.separator=: <<'EOF'
C001:{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
C002:{"order_id":"O-1002","customer_id":"C002","amount":500,"status":"CREATED"}
C001:{"order_id":"O-1003","customer_id":"C001","amount":88000,"status":"CANCELLED"}
EOF
연산하는 일입력 → 출력
filter조건에 맞는 것만 통과(C001,{...CREATED}),(C001,{...CANCELLED})CREATED
filterNot조건에 안 맞는 것만반대
map키와 값 둘 다 바꿈(C001, v)(O-1001, v)
mapValues값만 바꿈. 키 유지(C001, v)(C001, V)
flatMap하나를 0~N개로(C001, "a,b,c") → 3건
flatMapValues값만 0~N개로. 키 유지좌동
selectKey키만 바꿈(C001, v)(O-1001, v)
branch / split조건별로 여러 스트림으로 분기1스트림 → N스트림
merge여러 스트림을 하나로N스트림 → 1스트림
peek아무것도 안 바꾸고 들여다봄로깅·디버깅용
foreach종단. 다음 노드 없음사이드 이펙트용

출력 (앱 로그)

[PEEK-in ] C001 → {"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
[PEEK-in ] C002 → {"order_id":"O-1002","customer_id":"C002","amount":500,"status":"CREATED"}
[PEEK-in ] C001 → {"order_id":"O-1003","customer_id":"C001","amount":88000,"status":"CANCELLED"}
[branch  ] BIG   ← C001 amount=39000
[branch  ] SMALL ← C002 amount=500
[filter  ] CANCELLED 제외됨: O-1003
[merge   ] 총 2건이 s13_out 으로
docker exec kafka-1 /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_out \
  --from-beginning --property print.key=true --timeout-ms 5000

결과

C001	{"ORDER_ID":"O-1001","CUSTOMER_ID":"C001","AMOUNT":39000,"STATUS":"CREATED"}
C002	{"ORDER_ID":"O-1002","CUSTOMER_ID":"C002","AMOUNT":500,"STATUS":"CREATED"}
Processed a total of 2 messages

💡 mapmapValues 중에는 항상 mapValues 를 먼저 고려하세요 mapValues 는 키를 안 바꾼다고 Streams 에게 약속하는 것입니다. 그래서 파티션이 그대로 유지되고, 다음 절의 repartition 토픽이 안 생깁니다. map 은 키를 바꿀 수도 있다고 선언하는 것이라, 실제로 키를 안 바꿔도 Streams 는 "바뀌었을 수 있다"고 보고 repartition 을 겁니다. 같은 이유로 flatMap 보다 flatMapValues 를, transform 보다 transformValues 를 씁니다.


13-5. 함정 A 의 절반 — selectKey 뒤에 토픽이 생깁니다

지금 토픽 목록을 다시 봅니다.

kt --list

결과 (13-4 실행 후. 아직 안 늘었습니다)

__consumer_offsets
dlq
order-events
orders
payments
s13_left
s13_orders
s13_out
s13_right

스테이트리스 연산만으로는 토픽이 안 생깁니다. 이제 selectKey 를 넣고 그 뒤에 집계를 붙입니다.

builder.stream("s13_orders", Consumed.with(Serdes.String(), Serdes.String()))
       .selectKey((k, v) -> extractOrderId(v))   // customer_id → order_id 로 키 교체
       .groupByKey()
       .count(Materialized.as("order-counts"))
       .toStream()
       .to("s13_out", Produced.with(Serdes.String(), Serdes.Long()));

토폴로지를 봅니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java topology selectkey'

결과

Topologies:
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000000 (topics: [s13_orders])
      --> KSTREAM-KEY-SELECT-0000000001
    Processor: KSTREAM-KEY-SELECT-0000000001 (stores: [])
      --> order-counts-repartition-filter
      <-- KSTREAM-SOURCE-0000000000
    Processor: order-counts-repartition-filter (stores: [])
      --> order-counts-repartition-sink
      <-- KSTREAM-KEY-SELECT-0000000001
    Sink: order-counts-repartition-sink (topic: order-counts-repartition)
      <-- order-counts-repartition-filter

  Sub-topology: 1
    Source: order-counts-repartition-source (topics: [order-counts-repartition])
      --> KSTREAM-AGGREGATE-0000000002
    Processor: KSTREAM-AGGREGATE-0000000002 (stores: [order-counts])
      --> KTABLE-TOSTREAM-0000000006
      <-- order-counts-repartition-source
    Processor: KTABLE-TOSTREAM-0000000006 (stores: [])
      --> KSTREAM-SINK-0000000007
      <-- KSTREAM-AGGREGATE-0000000002
    Sink: KSTREAM-SINK-0000000007 (topic: s13_out)
      <-- KTABLE-TOSTREAM-0000000006

서브토폴로지가 둘로 쪼개졌습니다. 그 경계에 order-counts-repartition 토픽이 있습니다.

앱을 실행하고 토픽 목록을 다시 봅니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java selectkey' &
sleep 20
kt --list

결과

__consumer_offsets
dlq
order-events
orders
payments
s13-selectkey-app-order-counts-changelog
s13-selectkey-app-order-counts-repartition
s13_left
s13_orders
s13_out
s13_right

두 개가 새로 생겼습니다. 아무도 만들라고 하지 않았습니다.

왜 repartition 이 필요합니까

Streams 의 집계는 **"같은 키는 같은 태스크가 처리한다"**는 전제 위에 있습니다. 그래야 로컬 RocksDB 하나로 키별 상태를 관리할 수 있습니다.

   s13_orders (key=customer_id)
   ┌─────────┬─────────┬─────────┐
   │ P0      │ P1      │ P2      │
   │ C001    │ C002    │ C003    │
   └────┬────┴────┬────┴────┬────┘
        │ selectKey(order_id) ← 키가 바뀝니다
        ▼         ▼         ▼
   O-1001 이 P0 에, O-1002 도 P0 에, O-1003 은... 어디에?
   ★ 키는 바뀌었는데 레코드는 여전히 원래 파티션에 있습니다.
     같은 order_id 가 P0 과 P2 에 흩어져 있을 수 있습니다.

        │ repartition 토픽에 다시 씁니다 (key=order_id 로 파티셔닝)

   order-counts-repartition
   ┌─────────┬─────────┬─────────┐
   │ P0      │ P1      │ P2      │
   │ O-1003  │ O-1001  │ O-1002  │  ← 이제 같은 order_id 는 한 파티션에만
   └─────────┴─────────┴─────────┘

repartition 은 데이터를 Kafka 에 한 번 다시 쓰고 다시 읽는 것입니다. 공짜가 아닙니다. 처리량이 대략 절반이 되고 지연이 늘어납니다. 그래서 13-4 의 팁 — mapValues 를 쓰라 — 이 중요합니다.

repartition 토픽을 열어 봅니다.

kt --describe --topic s13-selectkey-app-order-counts-repartition

결과

Topic: s13-selectkey-app-order-counts-repartition	TopicId: mR8vK2nQTZ6yLpWcXbYhdA	PartitionCount: 3	ReplicationFactor: 1	Configs: cleanup.policy=delete,segment.bytes=52428800,retention.ms=-1,message.timestamp.type=CreateTime
	Topic: s13-selectkey-app-order-counts-repartition	Partition: 0	Leader: 1	Replicas: 1	Isr: 1
	Topic: s13-selectkey-app-order-counts-repartition	Partition: 1	Leader: 2	Replicas: 2	Isr: 2
	Topic: s13-selectkey-app-order-counts-repartition	Partition: 2	Leader: 3	Replicas: 3	Isr: 3

⚠️ 함정 — repartition 토픽의 RF 가 1 입니다 replication.factor 의 Streams 기본값은 1 입니다(3.x 기준). 브로커 기본값 default.replication.factor=3 을 따르지 않습니다. 즉 운영 클러스터에서 브로커 한 대가 죽으면 그 브로커가 리더였던 repartition/changelog 파티션이 통째로 사라지고, Streams 앱은 상태를 복구하지 못해 처음부터 다시 만들거나 아예 못 뜹니다. 해결: Streams 앱 설정에 반드시 넣으세요.

props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);

이 값은 repartition 토픽과 changelog 토픽에 모두 적용됩니다. 운영 체크리스트의 1번 항목입니다. 또 retention.ms=-1(무한)인 것도 눈여겨보세요. repartition 토픽은 Streams 가 소비 직후 스스로 purge 하므로 무한 보존이어도 안 쌓입니다. 이 삭제를 담당하는 것이 RepartitionTopics 의 purge 로직이며, repartition.purge.interval.ms(기본 30초)마다 돕니다.


13-6. 집계와 상태 저장소

집계 연산은 상태를 씁니다.

연산결과 타입하는 일
groupByKey()KGroupedStream키를 안 바꾸고 그룹핑. repartition 안 생김
groupBy((k,v) -> ...)KGroupedStream새 키로 그룹핑. repartition 생김
count()KTable<K, Long>키별 개수
reduce((a,b) -> ...)KTable<K, V>같은 타입끼리 접기
aggregate(초기값, (k,v,agg) -> ...)KTable<K, VR>다른 타입으로 접기
KTable<String, Long> counts =
    builder.stream("s13_orders", Consumed.with(Serdes.String(), Serdes.String()))
           .groupByKey()
           .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                  .withKeySerde(Serdes.String())
                  .withValueSerde(Serdes.Long()));
counts.toStream().to("s13_out", Produced.with(Serdes.String(), Serdes.Long()));
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java count' &
sleep 20
docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_orders \
  --property parse.key=true --property key.separator=: <<'EOF'
C001:{"order_id":"O-1001","amount":39000}
C002:{"order_id":"O-1002","amount":12500}
C001:{"order_id":"O-1003","amount":88000}
C001:{"order_id":"O-1004","amount":5000}
EOF
sleep 5
docker exec kafka-1 /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_out --from-beginning \
  --property print.key=true \
  --value-deserializer org.apache.kafka.common.serialization.LongDeserializer \
  --timeout-ms 5000

결과

C001	1
C002	1
C001	2
C001	3
Processed a total of 4 messages

C001 이 1, 2, 3 으로 세 번 나옵니다. KTable 은 "변경로그"이므로 값이 바뀔 때마다 갱신을 내보냅니다. "최종 결과 하나만" 원한다면 13-8 의 suppress 가 필요합니다.

상태 저장소 — RocksDB 와 changelog

상태는 두 곳에 있습니다.

   ① 로컬 RocksDB              ② Kafka changelog 토픽
   /tmp/kafka-streams/         s13-count-app-order-counts-changelog
     s13-count-app/
       0_0/rocksdb/order-counts/   ← 빠른 읽기용
       0_1/rocksdb/order-counts/       (프로세스가 죽으면 사라질 수 있음)
       0_2/rocksdb/order-counts/
                                  ← 진짜 원본. 재기동 시 여기서 복원

컨테이너 안에서 직접 봅니다.

docker exec kafka-1 find /tmp/kafka-streams -maxdepth 4 -type d | head -12

결과

/tmp/kafka-streams
/tmp/kafka-streams/s13-count-app
/tmp/kafka-streams/s13-count-app/0_0
/tmp/kafka-streams/s13-count-app/0_0/rocksdb
/tmp/kafka-streams/s13-count-app/0_0/rocksdb/order-counts
/tmp/kafka-streams/s13-count-app/0_1
/tmp/kafka-streams/s13-count-app/0_1/rocksdb
/tmp/kafka-streams/s13-count-app/0_1/rocksdb/order-counts
/tmp/kafka-streams/s13-count-app/0_2
/tmp/kafka-streams/s13-count-app/0_2/rocksdb
/tmp/kafka-streams/s13-count-app/0_2/rocksdb/order-counts
/tmp/kafka-streams/s13-count-app/0_0/rocksdb/order-counts/LOG

0_0, 0_1, 0_2태스크 ID 입니다. <서브토폴로지>_<파티션> 형식입니다. 입력 토픽이 3파티션이므로 태스크가 3개입니다.

changelog 토픽을 봅니다.

kt --describe --topic s13-count-app-order-counts-changelog

결과

Topic: s13-count-app-order-counts-changelog	TopicId: fW3pL7kNQR2mBcXvY8tZdg	PartitionCount: 3	ReplicationFactor: 1	Configs: cleanup.policy=compact,segment.bytes=52428800,message.timestamp.type=CreateTime
	Topic: s13-count-app-order-counts-changelog	Partition: 0	Leader: 2	Replicas: 2	Isr: 2
	Topic: s13-count-app-order-counts-changelog	Partition: 1	Leader: 3	Replicas: 3	Isr: 3
	Topic: s13-count-app-order-counts-changelog	Partition: 2	Leader: 1	Replicas: 1	Isr: 1

cleanup.policy=compact 입니다. repartition 토픽(delete)과 다릅니다.

repartitionchangelog
이름<app-id>-<노드/스토어>-repartition<app-id>-<스토어>-changelog
정책delete (+ Streams 가 스스로 purge)compact
담는 것키를 바꾼 레코드를 다시 흘려보냄상태 저장소의 모든 변경
언제 생김selectKey/map/groupBy/join 뒤에 집계·조인이 오면count/reduce/aggregate/table() 등 상태를 쓰면
지우면다시 만들어짐 (데이터는 재생성됨)상태를 잃습니다
파티션 수입력과 동일입력과 동일

changelog 가 compact 여야 하는 이유는 명확합니다. 키별 최신 값만 있으면 상태를 완전히 복원할 수 있기 때문입니다. C001 → 1, 2, 33 만 남아도 복원에 문제가 없습니다. delete 였다면 retention 이 지난 뒤 앱을 재기동했을 때 상태 일부가 사라진 채로 복원되고, 카운트가 조용히 작아집니다.

내용을 직접 봅니다.

docker exec kafka-1 /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka-1:9092 \
  --topic s13-count-app-order-counts-changelog --from-beginning \
  --property print.key=true \
  --value-deserializer org.apache.kafka.common.serialization.LongDeserializer \
  --timeout-ms 5000

결과

C001	1
C002	1
C001	2
C001	3
Processed a total of 4 messages

압축이 돌면 C001 3C002 1 만 남습니다.

💡 실무 팁 — 상태 복원 시간이 재기동 시간을 결정합니다 Streams 앱을 재기동하면 로컬 RocksDB 가 없는 경우(컨테이너 재생성, 새 인스턴스) changelog 토픽을 처음부터 끝까지 읽어 상태를 복원합니다. 상태가 수 GB 면 이 복원에 수십 분이 걸립니다. 그동안 앱은 데이터를 처리하지 않습니다. 해결 두 가지

  1. 상태 디렉터리를 영속 볼륨에 둡니다(state.dir). 그러면 재기동 시 로컬 것을 쓰고 델타만 따라잡습니다.
  2. num.standby.replicas=1 — 다른 인스턴스가 미리 상태 사본을 유지합니다. 장애 시 그 인스턴스가 즉시 인계받습니다. 디스크를 두 배 쓰는 대신 복구 시간이 거의 0 이 됩니다. 운영 Streams 앱에서 num.standby.replicas 는 사실상 필수 설정입니다.

13-7. 함정 A — 내부 토픽이 자동으로 생깁니다

앞의 두 절에서 이미 봤지만, 정면으로 다룹니다.

before (13-0 에서 기록한 것):

__consumer_offsets
dlq
order-events
orders
payments
s13_left
s13_orders
s13_out
s13_right

after (지금까지 앱 세 개를 띄운 뒤):

kt --list > /tmp/topics-after.txt
diff /tmp/topics-before.txt /tmp/topics-after.txt

결과

5a6,10
> s13-count-app-order-counts-changelog
> s13-selectkey-app-order-counts-changelog
> s13-selectkey-app-order-counts-repartition
> s13-window-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
> s13-window-app-KSTREAM-KEY-SELECT-0000000002-repartition

다섯 개가 자동으로 생겼습니다.

이름 규칙

패턴언제 생김
<app-id>-<스토어이름>-changelogMaterialized.as("이름") 으로 이름을 준 상태 저장소s13-count-app-order-counts-changelog
<app-id>-<노드이름>-changelog이름을 안 준 상태 저장소 (자동 생성 이름)s13-window-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
<app-id>-<스토어이름>-repartition이름 있는 스토어 앞의 키 변경s13-selectkey-app-order-counts-repartition
<app-id>-<노드이름>-repartition이름 없는 경우s13-window-app-KSTREAM-KEY-SELECT-0000000002-repartition
<app-id>-<이름>-subscription-registration-topicKTable 외래키 조인(이 코스에서는 안 나옴)

💡 실무 팁 — 상태 저장소에는 반드시 이름을 주세요 Materialized.as("order-counts") 를 쓰면 토픽 이름이 앱-order-counts-changelog 로 사람이 읽을 수 있게 나옵니다. 안 주면 앱-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog 가 됩니다. 문제는 뒤쪽 숫자가 DSL 호출 순서라는 점입니다. 코드에서 filter 를 하나 추가하면 번호가 밀리고, 토픽 이름이 바뀝니다. 새 이름의 토픽이 만들어지고 기존 상태는 고아가 됩니다. 배포 한 번에 집계가 리셋되는 것입니다. 이름을 주면 이 문제가 완전히 사라집니다. 운영 앱에서는 예외 없이 이름을 주세요.

auto.create.topics.enable=false 인데도 생기는 이유

우리 클러스터는 자동 토픽 생성이 꺼져 있습니다(Step 01, docker/index.md 참조).

kconf --describe --entity-type brokers --entity-name 1 | grep auto.create

결과

  auto.create.topics.enable=false sensitive=false synonyms={STATIC_BROKER_CONFIG:auto.create.topics.enable=false, DEFAULT_CONFIG:auto.create.topics.enable=true}

그런데 Streams 는 만들었습니다.

⚠️ 함정 A — Streams 는 AdminClient 로 토픽을 직접 만듭니다 auto.create.topics.enable 은 **"프로듀서/컨슈머가 없는 토픽에 접근하면 자동 생성할까"**를 결정하는 설정입니다. Streams 는 그 경로를 안 씁니다. 기동 시 InternalTopicManagerAdminClient 로 CreateTopics API 를 명시적으로 호출합니다. 사용자가 손으로 kt --create 하는 것과 완전히 같은 경로입니다. 그래서 자동 생성을 꺼도 막히지 않습니다. Connect(Step 12)의 내부 토픽도 같은 방식입니다.

왜 문제입니까 운영 클러스터에서 "누가 이 토픽 만들었지?"의 정체가 대개 이것입니다. 토픽 목록이 어느 날부터 두 배가 되어 있고, 이름은 아무도 모르는 app-KSTREAM-AGGREGATE-STATE-STORE-0000000019-changelog 같은 것입니다. 더 나쁜 것은 앱을 지워도 토픽은 남는다는 점입니다. Streams 앱을 내리고 코드를 지워도, 토픽·컨슈머 그룹·상태 디렉터리가 전부 남습니다. 디스크를 계속 먹습니다.

해결

  1. application.id명확한 접두사를 씁니다. s13-count-app 처럼요. 내부 토픽이 전부 그 접두사로 시작하므로 소유자를 즉시 알 수 있습니다.
  2. 상태 저장소에 이름을 줍니다. 위 팁 참조.
  3. 앱을 폐기할 때는 반드시 kafka-streams-application-reset.sh 를 돌립니다. 다음 절입니다.
  4. RF 를 3 으로 (StreamsConfig.REPLICATION_FACTOR_CONFIG). 기본값 1 은 운영에서 위험합니다.

kafka-streams-application-reset.sh

앱의 흔적을 정리하는 공식 도구입니다. 앱이 완전히 멈춘 상태여야 합니다.

docker exec kafka-1 /opt/kafka/bin/kafka-streams-application-reset.sh \
  --bootstrap-server kafka-1:9092 \
  --application-id s13-selectkey-app \
  --input-topics s13_orders

결과

Reset-offsets for input topics [s13_orders]
Following input topics offsets will be reset to (for consumer group s13-selectkey-app)
Topic: s13_orders Partition: 0 Offset: 0
Topic: s13_orders Partition: 1 Offset: 0
Topic: s13_orders Partition: 2 Offset: 0
Done.
Deleting all internal/auto-created topics for application s13-selectkey-app
Deleted topic: s13-selectkey-app-order-counts-changelog
Deleted topic: s13-selectkey-app-order-counts-repartition
Done.

앱이 아직 돌고 있으면 이렇게 거부합니다.

ERROR: Java class 'kafka.tools.StreamsResetter' failed:
  java.lang.IllegalStateException: Consumer group 's13-selectkey-app' is still active
  and has following members: [s13-selectkey-app-8f2c...-StreamThread-1-consumer].
  Make sure to stop all running application instances before running the reset tool.

주요 옵션:

옵션하는 일
--application-id필수. 정리할 앱
--input-topics입력 토픽의 오프셋을 처음으로 되돌림
--intermediate-topicsthrough() 로 쓴 중간 토픽을 비움
--to-datetime / --to-offset / --shift-by처음이 아닌 특정 지점으로
--dry-run실제로 안 하고 계획만 출력
--force활성 멤버가 있어도 강제 (그룹에서 쫓아냄)

⚠️ 함정 — reset 도구가 로컬 상태 디렉터리는 안 지웁니다 이 도구는 Kafka 쪽만 정리합니다. 각 인스턴스의 state.dir(기본 /tmp/kafka-streams/<app-id>)에 남은 RocksDB 파일은 그대로입니다. 오프셋만 0 으로 되돌리고 로컬 상태는 남은 채로 앱을 재기동하면, 옛 상태에 새로 읽은 데이터가 얹혀서 카운트가 두 배가 됩니다. 에러는 없습니다. 해결: 각 인스턴스에서 KafkaStreams.cleanUp() 을 기동 전에 호출하거나, 디렉터리를 직접 지웁니다.

docker exec kafka-1 rm -rf /tmp/kafka-streams/s13-selectkey-app

cleanUp()start() 전에만 호출할 수 있습니다. 운영 앱에는 --reset 같은 기동 플래그를 만들어 두고, 그때만 cleanUp() 을 부르게 하는 패턴이 흔합니다. 무조건 부르면 재기동마다 전체 복원이 일어나 기동이 몇십 분씩 걸립니다.


13-8. 윈도우

시간 구간으로 잘라 집계합니다.

종류정의겹침한 레코드가 속하는 윈도우 수용도
Tumbling고정 크기, 겹치지 않음1"1분당 주문 수"
Hopping고정 크기 + advance 간격size / advance"최근 5분, 1분마다 갱신"
Sliding레코드 기준 ± 크기가변"이 이벤트 앞뒤 10초 안의 것들"
Session비활동 간격(gap)으로 구분1"사용자 세션당 행동 수"
// Tumbling — 1분
TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30))

// Hopping — 5분 크기, 1분마다 (한 레코드가 5개 윈도우에 들어감)
TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30))
           .advanceBy(Duration.ofMinutes(1))

// Sliding — 10초
SlidingWindows.ofTimeDifferenceAndGrace(Duration.ofSeconds(10), Duration.ofSeconds(5))

// Session — 5분 비활동
SessionWindows.ofInactivityGapAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30))

1분 텀블링 — 고객별 주문 금액 합계

builder.stream("s13_orders", Consumed.with(Serdes.String(), Serdes.String()))
       .groupByKey()
       .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
       .aggregate(() -> 0L,
                  (k, v, agg) -> agg + amountOf(v),
                  Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("amount-per-min")
                              .withKeySerde(Serdes.String())
                              .withValueSerde(Serdes.Long()))
       .toStream()
       .foreach((wk, sum) -> System.out.printf("[WINDOW] %s [%s ~ %s] sum=%d%n",
                wk.key(), wk.window().startTime(), wk.window().endTime(), sum));
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java window' &
sleep 20
docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_orders \
  --property parse.key=true --property key.separator=: <<'EOF'
C001:{"order_id":"O-1001","amount":10000}
C001:{"order_id":"O-1002","amount":20000}
C002:{"order_id":"O-1003","amount":5000}
EOF

결과 (앱 로그)

[WINDOW] C001 [2024-03-11T10:22:00Z ~ 2024-03-11T10:23:00Z] sum=10000
[WINDOW] C001 [2024-03-11T10:22:00Z ~ 2024-03-11T10:23:00Z] sum=30000
[WINDOW] C002 [2024-03-11T10:22:00Z ~ 2024-03-11T10:23:00Z] sum=5000

윈도우 경계가 10:22:00 ~ 10:23:00 으로 깔끔합니다. 텀블링 윈도우는 에폭(1970-01-01T00:00:00Z)부터 크기 단위로 잘린 절대 시각을 씁니다. 앱이 언제 시작했든 같은 경계가 나옵니다.

함정 B — grace period 를 지난 레코드는 조용히 버려집니다

Streams 의 윈도우는 이벤트 시간(레코드의 타임스탬프)으로 동작합니다. 네트워크 지연이나 재시도 때문에 레코드가 늦게 도착할 수 있는데, 언제까지 기다릴지가 grace 입니다.

일부러 늦은 레코드를 만듭니다. 콘솔 프로듀서로는 타임스탬프를 못 지정하므로 Practice.javawindow 시나리오가 대신 만들어 줍니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java window --inject-late'

결과

[PRODUCE] C001 amount=10000  ts=2024-03-11T10:22:10Z  (윈도우 10:22:00~10:23:00, 정상)
[PRODUCE] C001 amount=20000  ts=2024-03-11T10:22:40Z  (같은 윈도우, 정상)
[WINDOW ] C001 [10:22:00 ~ 10:23:00] sum=10000
[WINDOW ] C001 [10:22:00 ~ 10:23:00] sum=30000

[PRODUCE] C001 amount=7000   ts=2024-03-11T10:22:50Z  (스트림 시간은 이미 10:23:20, grace 30초 → 10:23:30 까지 허용. 아직 유효)
[WINDOW ] C001 [10:22:00 ~ 10:23:00] sum=37000

[PRODUCE] C001 amount=99000  ts=2024-03-11T10:22:55Z  (스트림 시간 10:23:45. grace 만료됨)
   ★ 아무 출력도 없습니다. 예외도 없습니다.

[FINAL  ] C001 [10:22:00 ~ 10:23:00] sum=37000     ← 99000 이 빠졌습니다
[METRIC ] dropped-records-total = 1

⚠️ 함정 B — 늦게 도착한 레코드는 예외 없이, 로그 없이 사라집니다 TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)) 은 "윈도우가 끝난 뒤 30초까지만 늦은 레코드를 받는다"는 뜻입니다. 31초째에 온 레코드는 버려집니다. 예외도 안 던지고, 기본 로그 레벨에서는 아무것도 안 찍힙니다. 결과 토픽의 합계가 그냥 작습니다. "Kafka 집계 결과가 DB 집계와 안 맞는다"의 가장 흔한 원인입니다.

어떻게 알아챕니까 — 지표를 보세요.

kafka.streams:type=stream-task-metrics,thread-id=...,task-id=1_0
  dropped-records-total   ← 이 값이 0 이 아니면 데이터가 버려지고 있습니다
  dropped-records-rate

이 지표에 알람을 거는 것이 유일한 방어입니다. Step 14 의 JMX 절과 이어집니다.

grace 를 얼마로 잡습니까

  • 너무 짧으면: 정상 데이터가 버려집니다.
  • 너무 길면: 최종 결과가 그만큼 늦게 나오고, 윈도우 상태를 그동안 메모리/디스크에 들고 있어야 합니다.
  • 실무 기준: 관측된 최대 지연의 1.5~2배. 프로듀서 재시도(delivery.timeout.ms)와 네트워크 지연을 합쳐서 계산합니다. 처음에는 넉넉히 잡고 dropped-records-total 을 보며 줄이는 편이 안전합니다.

⚠️ 3.0 에서 기본값이 바뀌었습니다. TimeWindows.of(size) 는 deprecated 되었고, 옛 기본 grace 는 24시간이었습니다. 새 ofSizeAndGrace(size, grace) 는 grace 를 명시하도록 강제합니다. 옛 코드를 마이그레이션할 때 grace 를 안 주면 ofSizeWithNoGrace() 로 바뀌어 grace 0 이 되고, 조금만 늦어도 전부 버려집니다. 3.x 마이그레이션의 대표적인 사고 지점입니다.

suppress — 최종 결과만 내보내기

위 출력에서 sum=10000, sum=30000, sum=37000 이 순차적으로 나왔습니다. 다운스트림이 DB 업서트라면 세 번 쓰게 됩니다. 윈도우가 닫힌 뒤 최종값 하나만 원한다면:

.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
.aggregate(...)
.suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
.toStream()
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java suppress'

결과

(윈도우 진행 중 — 아무 출력 없음)
(10:23:30 — grace 만료, 윈도우 닫힘)
[FINAL] C001 [10:22:00 ~ 10:23:00] sum=37000
[FINAL] C002 [10:22:00 ~ 10:23:00] sum=5000

중간 결과가 사라지고 최종 하나씩만 나옵니다.

⚠️ 함정 — suppress 는 다음 레코드가 와야 방출합니다 untilWindowCloses 는 "스트림 시간이 윈도우 끝 + grace 를 넘으면 방출"입니다. 그런데 스트림 시간은 레코드가 들어와야 전진합니다. 트래픽이 멈추면 스트림 시간도 멈추고, 마지막 윈도우는 영원히 방출되지 않습니다. 개발 중에 "suppress 를 걸었더니 아무것도 안 나온다"는 열에 아홉 이것입니다. 벽시계로 30초를 기다려도 안 나옵니다. 레코드를 하나 더 넣으면 그제서야 나옵니다. 해결: 저트래픽 토픽에는 하트비트 레코드를 주기적으로 넣거나, suppress 대신 중간 결과를 받아 다운스트림에서 멱등 업서트로 처리합니다. 후자가 더 견고합니다.

BufferConfig.unbounded() 도 위험합니다. 방출 전까지 모든 윈도우를 메모리에 들고 있습니다. 키가 많으면 OOM 입니다. 운영에서는 상한을 주세요.

Suppressed.BufferConfig.maxBytes(50_000_000L).emitEarlyWhenFull()

emitEarlyWhenFull() 은 버퍼가 차면 조기 방출합니다(= suppress 효과 일부 포기). 대안인 shutDownWhenFull() 은 앱을 죽입니다. 조용히 OOM 나는 것보다는 낫습니다만, 대개 emitEarlyWhenFull 이 실용적입니다.


13-9. 조인

orders × payments 윈도우 조인

KStream<String, String> orders   = builder.stream("s13_orders");
KStream<String, String> payments = builder.stream("s13_payments");

orders.selectKey((k, v) -> orderIdOf(v))          // customer_id → order_id
      .join(payments,
            (o, p) -> "{\"order\":" + o + ",\"payment\":" + p + "}",
            JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30)),
            StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String()))
      .to("s13_out");
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java join' &
sleep 20
docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_orders \
  --property parse.key=true --property key.separator=: <<'EOF'
C001:{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
EOF
sleep 2
docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_payments \
  --property parse.key=true --property key.separator=: <<'EOF'
O-1001:{"order_id":"O-1001","method":"CARD","amount":39000,"result":"APPROVED"}
EOF
sleep 3
docker exec kafka-1 /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka-1:9092 --topic s13_out --from-beginning \
  --property print.key=true --timeout-ms 5000

결과

O-1001	{"order":{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"},"payment":{"order_id":"O-1001","method":"CARD","amount":39000,"result":"APPROVED"}}
Processed a total of 1 messages

조인 종류별 특성

조합윈도우상태 저장소트리거
KStream-KStream필수양쪽 다 (윈도우 스토어 2개)어느 쪽이 와도
KStream-KTable없음KTable 쪽만스트림 쪽이 올 때만
KTable-KTable없음양쪽 다어느 쪽이 바뀌어도
KStream-GlobalKTable없음GlobalKTable (전체 복제)스트림 쪽이 올 때만

KStream-KTable 조인은 스트림 쪽이 올 때만 발화합니다. KTable 이 나중에 갱신돼도 이미 지나간 스트림 레코드를 다시 조인하지 않습니다. "주문이 왔는데 고객 정보가 아직 KTable 에 없어서 조인이 안 됐다"는 상황이 흔하고, 나중에 고객 정보가 와도 그 주문은 영영 조인되지 않습니다.

함정 C — co-partitioning

s13_left 는 3파티션, s13_right 는 6파티션입니다(13-0 에서 그렇게 만들었습니다). 조인해 봅니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java copartition'

결과

Exception in thread "main" org.apache.kafka.streams.errors.TopologyException: Invalid topology: Following topics do not have the same number of partitions: [s13_left(3), s13_right(6)]
	at org.apache.kafka.streams.processor.internals.InternalTopologyBuilder.verifyCopartitioning(InternalTopologyBuilder.java:1149)
	at org.apache.kafka.streams.processor.internals.StreamsPartitionAssignor.assign(StreamsPartitionAssignor.java:406)
	at Practice$CoPartition.run(Practice.java:412)
	at Practice.main(Practice.java:58)

⚠️ 함정 C — 조인하려면 co-partitioning 이 필요합니다 Streams 의 조인은 **"같은 키는 같은 태스크가 본다"**는 전제 위에 있습니다. 태스크 0_1 은 양쪽 토픽의 파티션 1 만 봅니다. 왼쪽의 O-1001 이 파티션 1 에 있는데 오른쪽의 O-1001 이 파티션 4 에 있으면, 어느 태스크도 두 레코드를 함께 볼 수 없습니다.

co-partitioning 의 세 조건

  1. 파티션 수가 같아야 합니다. — 위 예외가 이것입니다. 유일하게 Streams 가 검증해 주는 조건입니다.
  2. 키의 타입과 직렬화가 같아야 합니다. — 한쪽이 String "1001", 다른 쪽이 Long 1001L 이면 바이트가 달라 다른 파티션으로 갑니다.
  3. 파티셔너가 같아야 합니다. — ★ 이것은 검증되지 않습니다.

3번이 진짜 함정입니다. 예를 들어 한쪽 토픽을 커스텀 파티셔너를 쓰는 Java 프로듀서가 채우고, 다른 쪽을 기본 파티셔너로 채웠다면, 파티션 수가 같아도 같은 키가 다른 파티션에 들어갑니다. Streams 는 아무 예외도 안 냅니다. 조인이 그냥 안 됩니다. 결과 토픽이 비어 있거나 일부만 나옵니다. 코드를 아무리 봐도 이상한 곳이 없습니다.

해결

  • 1번(파티션 수)은 조인 전에 한쪽을 repartition(Repartitioned.numberOfPartitions(3)) 으로 맞춥니다.
    KStream<String,String> right = builder.stream("s13_right")
        .repartition(Repartitioned.<String,String>as("right-fixed").numberOfPartitions(3));
    s13-copart-app-right-fixed-repartition 이라는 3파티션 토픽이 생기고, 그것과 조인합니다.
  • 3번(파티셔너)은 양쪽을 같은 방식으로 쓰는 것이 유일한 예방책입니다. Kafka 기본 파티셔너(murmur2)를 쓰거나, 정 안 되면 조인 전에 양쪽 다 repartition() 을 태워 Streams 의 파티셔너로 통일합니다.
  • GlobalKTable 은 이 문제를 원천적으로 회피합니다. 전체를 복제하므로 파티셔닝이 무의미합니다. 작은 참조 데이터라면 이게 가장 편합니다.

수정한 버전으로 다시 실행합니다.

docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java copartition --fix'

결과

[FIX] s13_right(6) → repartition(3) 로 맞춥니다
[TOPOLOGY]
   Sub-topology: 0
    Source: KSTREAM-SOURCE-0000000001 (topics: [s13_right])
      --> right-fixed-repartition-filter
    ...
    Sink: right-fixed-repartition-sink (topic: right-fixed-repartition)

  Sub-topology: 1
    Source: KSTREAM-SOURCE-0000000000 (topics: [s13_left])
    Source: right-fixed-repartition-source (topics: [right-fixed-repartition])
    ...
[JOIN] K-1  ← left="L1" right="R1"
[JOIN] K-2  ← left="L2" right="R2"
정상 동작합니다.

13-10. EOS — exactly_once_v2

Step 07 에서 트랜잭션 API 로 exactly-once 를 직접 구현했습니다. Streams 에서는 설정 한 줄입니다.

props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java eos' &
sleep 20

결과 (앱 로그)

[CONFIG] processing.guarantee = exactly_once_v2
[CONFIG] → 내부적으로 강제되는 값:
           enable.idempotence          = true
           max.in.flight.requests...   = 5
           acks                        = all
           transactional.id            = s13-eos-app-<uuid>-<taskId>
           isolation.level             = read_committed
           commit.interval.ms          = 100   (기본 30000 에서 변경됨)
[TXN] begin  → process 3 records → sendOffsetsToTransaction → commit

Streams 의 EOS 가 보장하는 것: 입력 소비 오프셋 커밋 + 상태 저장소 갱신 + 출력 토픽 쓰기, 이 셋이 하나의 트랜잭션으로 묶입니다. 셋 중 하나라도 실패하면 전부 롤백됩니다.

at_least_once (기본)exactly_once_v2
중복재시작 시 발생 가능없음
상태 일관성재시작 시 상태와 오프셋이 어긋날 수 있음항상 일치
commit.interval.ms 기본30000100
지연낮음커밋 간격만큼 추가 (기본 100ms)
처리량기준대략 -10~20%
컨슈머 요구없음다운스트림이 read_committed 여야

exactly_once_v2 는 Kafka 2.5+ 부터이며 브로커도 2.5 이상이어야 합니다. 옛 exactly_once(v1)는 태스크마다 프로듀서를 만들어 자원을 크게 썼는데, v2 는 스레드당 프로듀서 하나로 줄여 태스크가 많아도 확장됩니다. 3.0 부터 v1 은 deprecated 이고 4.0 에서 제거되었습니다. 새로 쓰면 무조건 v2 입니다.

⚠️ 함정 — EOS 는 Kafka 안에서만 유효합니다 Streams 의 EOS 는 "Kafka → 처리 → Kafka" 경로만 보장합니다. 처리 로직 안에서 외부 시스템을 호출하면(REST API, DB INSERT, 이메일 발송) 그것은 트랜잭션 밖입니다. 트랜잭션이 롤백되면 Kafka 쪽 쓰기는 취소되지만, 이미 보낸 이메일은 취소되지 않습니다. 재처리 시 이메일이 두 번 갑니다. 해결: 외부 호출은 멱등하게 만들거나(업서트, 멱등 키), Streams 로는 Kafka 토픽까지만 쓰고 외부 반영은 싱크 커넥터에게 맡깁니다(Step 12 의 insert.mode=upsert). 후자가 아키텍처적으로 깔끔합니다.


13-11. 정리 (실습 마무리)

Streams 실습은 정리할 것이 셋입니다. 토픽, 컨슈머 그룹, 로컬 상태 디렉터리.

먼저 실행 중인 앱을 전부 내립니다.

docker exec kafka-1 pkill -f 'Practice' || true
sleep 5
kcg --list

결과

s13-count-app
s13-eos-app
s13-join-app
s13-selectkey-app
s13-stateless-app
s13-window-app

각 앱을 reset 도구로 정리합니다.

for app in s13-stateless-app s13-selectkey-app s13-count-app s13-window-app s13-join-app s13-eos-app; do
  docker exec kafka-1 /opt/kafka/bin/kafka-streams-application-reset.sh \
    --bootstrap-server kafka-1:9092 --application-id "$app" \
    --input-topics s13_orders 2>&1 | tail -2
done

결과

Deleted topic: s13-stateless-app-... (없으면 아무것도 안 나옴)
Done.
Deleted topic: s13-selectkey-app-order-counts-changelog
Deleted topic: s13-selectkey-app-order-counts-repartition
Done.
Deleted topic: s13-count-app-order-counts-changelog
Done.
Deleted topic: s13-window-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
Deleted topic: s13-window-app-KSTREAM-KEY-SELECT-0000000002-repartition
Done.
Deleted topic: s13-join-app-KSTREAM-JOINTHIS-0000000009-store-changelog
Deleted topic: s13-join-app-KSTREAM-JOINOTHER-0000000010-store-changelog
Deleted topic: s13-join-app-KSTREAM-KEY-SELECT-0000000002-repartition
Done.
Deleted topic: s13-eos-app-order-counts-changelog
Done.

reset 도구는 컨슈머 그룹을 안 지웁니다. 오프셋만 0 으로 되돌립니다. 그룹도 지웁니다.

for g in $(kcg --list | grep '^s13-'); do kcg --delete --group "$g"; done
kcg --list

결과

Deletion of requested consumer groups ('s13-stateless-app') was successful.
Deletion of requested consumer groups ('s13-selectkey-app') was successful.
Deletion of requested consumer groups ('s13-count-app') was successful.
Deletion of requested consumer groups ('s13-window-app') was successful.
Deletion of requested consumer groups ('s13-join-app') was successful.
Deletion of requested consumer groups ('s13-eos-app') was successful.

로컬 상태 디렉터리를 지웁니다. 이걸 안 하면 다음에 같은 application.id 로 띄웠을 때 옛 상태가 되살아납니다.

docker exec kafka-1 du -sh /tmp/kafka-streams 2>/dev/null || echo "없음"
docker exec kafka-1 rm -rf /tmp/kafka-streams
docker exec kafka-1 ls /tmp/kafka-streams 2>&1 || echo "삭제 완료"

결과

3.4M	/tmp/kafka-streams
ls: cannot access '/tmp/kafka-streams': No such file or directory
삭제 완료

실습 토픽을 지웁니다.

kt --list | grep -E '^s13' | xargs -I{} docker exec kafka-1 \
  /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka-1:9092 --delete --topic {}
kt --list | grep -E '^s13' || echo "s13 토픽 없음"

결과

s13 토픽 없음
docker exec kafka-1 rm -f /tmp/Practice.java /tmp/topics-before.txt /tmp/topics-after.txt

정리

개념핵심
Streams 의 정체라이브러리. Consumer + Producer + 로컬 상태. 별도 클러스터 없음
확장같은 application.id 로 프로세스를 더 띄우면 끝
application.id= 컨슈머 그룹 ID. 랙 모니터링이 일반 컨슈머와 동일
Topology서브토폴로지 개수 = Kafka 네트워크 왕복 횟수
KStream레코드 스트림. 사건의 나열
KTable변경로그 스트림. 키별 현재 상태. null = 삭제
GlobalKTable전체 복제. co-partition 불필요. 작은 참조 데이터만
map vs mapValuesmap 은 키를 안 바꿔도 repartition 을 유발
함정 A내부 토픽이 자동 생성됨. auto.create.topics.enable=false 여도 AdminClient 로 직접 만듦
이름 규칙<app-id>-<이름>-repartition / <app-id>-<이름>-changelog
repartitiondelete 정책. Streams 가 스스로 purge. RF 기본 1
changelogcompact 정책. 상태의 진짜 원본. 지우면 상태를 잃음
상태 저장소로컬 RocksDB + changelog 토픽 이중화
reset 도구Kafka 쪽만 정리. 로컬 state.dir 은 따로 지워야 함
윈도우 4종Tumbling / Hopping / Sliding / Session
함정 Bgrace 지난 레코드는 예외·로그 없이 버려짐. dropped-records-total 로만 감지
suppress최종 결과만. 다음 레코드가 와야 방출됨 (저트래픽에서 안 나옴)
함정 C조인은 co-partitioning 필요. 파티션 수만 검증되고 파티셔너 불일치는 조용히 실패
EOSexactly_once_v2 한 줄. Kafka 안에서만 유효. 외부 호출은 별도 멱등화

연습문제

exercise.sh 에 6문제가 있습니다. 정답은 solution.sh.

  1. selectKey 를 쓴 앱과 안 쓴 앱의 kt --list 차이를 확인하고, 생긴 토픽 이름을 설명하기
  2. 어떤 앱이 만든 내부 토픽인지 이름만 보고 판별하기 (5개 중 소유 앱 맞히기)
  3. changelog 토픽의 cleanup.policy 를 확인하고, delete 였다면 무슨 일이 생기는지 답하기
  4. kafka-streams-application-reset.sh 로 앱을 완전히 정리하기 (로컬 상태 포함)
  5. 파티션 수가 다른 두 토픽을 조인해 TopologyException 을 재현하고 고치기
  6. 윈도우 집계 결과가 실제보다 작습니다. 어느 지표를 보고 어떻게 고칩니까?

다음 단계

Connect 로 데이터를 들여오고, Streams 로 가공했습니다. 이제 이 모든 것이 운영에서 계속 돌아가게 만들 차례입니다.

파티션을 옮기고, 브로커를 무중단으로 재시작하고, JMX 지표에 알람을 걸고, 장애가 나면 플레이북을 따라 대응합니다. 마지막으로 Step 01~13 을 전부 동원해 주문 이벤트 파이프라인을 처음부터 끝까지 구축하고, 브로커를 죽여 가며 유실 0 을 검증합니다.

Step 14 — 운영과 최종 프로젝트


실습 파일

이 스텝은 파일 네 개로 진행합니다. 중심은 Practice.java 이고, 셸 스크립트 셋은 그것을 띄우고 관찰하고 정리하는 역할입니다. 순서는 practice.shexercise.shsolution.sh 이며, Practice.java 는 세 스크립트가 모두 docker cp 로 컨테이너에 넣어 씁니다.

Practice.java

시나리오 9종을 인자로 골라 실행하는 단일 파일 Java 프로그램입니다. Java 21 의 single-file source 실행(java Practice.java <arg>)을 쓰므로 컴파일 단계가 없고, 의존성은 /opt/kafka/libs/* 뿐입니다.

  • 각 시나리오는 Scenario 인터페이스를 구현한 static 중첩 클래스입니다. Stateless, SelectKey, Count, Window, Suppress, Join, CoPartition, Eos, TopologyOnly 아홉 개이고, mainswitch 가 인자로 고릅니다.
  • TopologyOnly앱을 실행하지 않고 topology.describe() 만 출력하고 끝납니다. 13-2 와 13-5 의 서브토폴로지 출력을 재현할 때 씁니다. Practice.java topology selectkey 처럼 두 번째 인자로 어느 토폴로지를 볼지 고릅니다.
  • 모든 시나리오가 공통 baseProps(appId) 를 씁니다. 여기에 REPLICATION_FACTOR_CONFIG=3 이 들어 있습니다. Streams 기본값은 1 이라 그대로 두면 브로커 한 대만 죽어도 내부 토픽이 사라지는데, 그것을 막는 설정입니다. 13-5 의 함정 블록에서 설명한 그 한 줄입니다.
  • Window 시나리오는 --inject-late 플래그를 받습니다. 이 플래그가 있으면 프로듀서가 타임스탬프를 직접 지정해서 grace 를 넘긴 레코드를 하나 넣습니다. 콘솔 프로듀서로는 타임스탬프를 못 정하기 때문에 Java 로만 재현할 수 있는 실습입니다. 실행 끝에 dropped-records-total 지표를 읽어 출력하므로, 1 이 찍히는 것이 함정 B 의 증거입니다.
  • CoPartition 은 인자 없이 실행하면 s13_left(3) × s13_right(6) 조인을 시도해 TopologyException 을 그대로 터뜨립니다. --fix 를 주면 repartition(Repartitioned.numberOfPartitions(3)) 로 맞춘 버전이 돌아가며, 두 경우의 토폴로지 출력을 나란히 볼 수 있게 describe() 를 먼저 찍습니다.
  • Eos 는 실행 시 processing.guarantee=exactly_once_v2강제로 바꾸는 설정 5종(idempotence, acks, transactional.id, isolation.level, commit.interval.ms)을 로그로 출력합니다. Step 07 에서 손으로 설정했던 것들이 한 줄로 대체되는 것을 확인하는 용도입니다.
  • 모든 시나리오에 Runtime.getRuntime().addShutdownHook(...) 으로 streams.close(Duration.ofSeconds(10)) 가 걸려 있습니다. Ctrl+Cpkill 로 죽여도 컨슈머 그룹에서 깨끗이 빠져나가므로, 13-11 의 reset 도구가 "still active" 로 거부하는 일이 줄어듭니다.
// ============================================================================
// Step 13 — Kafka Streams / Practice.java
//
// 실행법 (Java 21 single-file source 실행. 컴파일 단계 없음):
//   docker cp Practice.java kafka-1:/tmp/
//   docker exec kafka-1 sh -c 'cd /tmp && java -cp "/opt/kafka/libs/*" Practice.java <scenario>'
//
// 의존성은 Kafka 배포판의 /opt/kafka/libs/* 뿐입니다.
//   kafka-streams-3.7.1.jar, kafka-clients-3.7.1.jar, rocksdbjni-*.jar 등이 들어 있습니다.
//
// 시나리오:
//   topology [name]  토폴로지만 출력하고 종료 (name: stateless|selectkey|count|window|join)
//   stateless        filter / map / mapValues / flatMap / branch / merge / peek
//   selectkey        selectKey → repartition 토픽 생성 관찰
//   count            groupByKey().count() → changelog 토픽과 RocksDB
//   window           1분 텀블링 윈도우. --inject-late 로 grace 초과 레코드 주입
//   suppress         suppress(untilWindowCloses) 로 최종 결과만
//   join             orders × payments 윈도우 조인
//   copartition      co-partitioning 위반 재현. --fix 로 수정판
//   eos              processing.guarantee=exactly_once_v2
//
// 사전 토픽 (practice.sh 가 만듭니다):
//   s13_orders(3) s13_payments(3) s13_out(3) s13_left(3) s13_right(6)
// ============================================================================

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.*;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.*;

import java.time.Duration;
import java.time.Instant;
import java.util.*;
import java.util.concurrent.CountDownLatch;

public class Practice {

    static final String BOOTSTRAP = "kafka-1:9092";

    static final String T_ORDERS   = "s13_orders";
    static final String T_PAYMENTS = "s13_payments";
    static final String T_OUT      = "s13_out";
    static final String T_LEFT     = "s13_left";
    static final String T_RIGHT    = "s13_right";

    public static void main(String[] args) throws Exception {
        if (args.length == 0) { usage(); return; }
        String scenario = args[0];
        String opt = args.length > 1 ? args[1] : "";

        switch (scenario) {
            case "topology"    -> TopologyOnly.run(opt);
            case "stateless"   -> new Stateless().run(opt);
            case "selectkey"   -> new SelectKey().run(opt);
            case "count"       -> new Count().run(opt);
            case "window"      -> new Window().run(opt);
            case "suppress"    -> new Suppress().run(opt);
            case "join"        -> new Join().run(opt);
            case "copartition" -> new CoPartition().run(opt);
            case "eos"         -> new Eos().run(opt);
            default            -> usage();
        }
    }

    static void usage() {
        System.out.println("""
            usage: java -cp "/opt/kafka/libs/*" Practice.java <scenario>

              topology   토폴로지 설명만 출력하고 종료 (실행하지 않음)
              stateless  filter / map / mapValues / flatMap / branch / merge / peek
              selectkey  selectKey → repartition 토픽이 생기는 것을 관찰
              count      groupByKey().count() → changelog 토픽과 상태 저장소
              window     1분 텀블링 윈도우 집계 + grace period
              suppress   suppress(untilWindowCloses) 로 최종 결과만
              join       orders × payments 윈도우 조인
              copartition co-partitioning 위반 재현 (TopologyException)
              eos        processing.guarantee=exactly_once_v2

              예) java -cp "/opt/kafka/libs/*" Practice.java count
            """);
    }

    // ------------------------------------------------------------------------
    // 공통 설정
    //
    // ★ REPLICATION_FACTOR_CONFIG = 3 이 중요합니다.
    //   Streams 의 기본값은 1 입니다. 브로커 기본값(default.replication.factor=3)을
    //   따르지 않습니다. 그대로 두면 repartition/changelog 토픽이 RF=1 로 만들어지고,
    //   브로커 한 대가 죽는 순간 상태를 잃습니다. 운영 체크리스트 1번 항목입니다.
    // ------------------------------------------------------------------------
    static Properties baseProps(String appId) {
        Properties p = new Properties();
        p.put(StreamsConfig.APPLICATION_ID_CONFIG, appId);
        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());
        p.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);        // ★ 기본 1 → 3
        p.put(StreamsConfig.STATE_DIR_CONFIG, "/tmp/kafka-streams");
        p.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1);
        p.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);     // 관찰을 빠르게 (기본 30000)
        p.put(StreamsConfig.consumerPrefix("auto.offset.reset"), "earliest");
        return p;
    }

    // 값 JSON 에서 필드 하나를 아주 단순하게 뽑습니다.
    // 실무라면 Jackson 을 쓰지만, 여기서는 /opt/kafka/libs/* 밖의 의존성을 안 쓰기 위해
    // 문자열 파싱으로 대신합니다. 값 포맷이 고정된 실습이라 문제없습니다.
    static String field(String json, String name) {
        String needle = "\"" + name + "\":";
        int i = json.indexOf(needle);
        if (i < 0) return null;
        int s = i + needle.length();
        while (s < json.length() && (json.charAt(s) == ' ' || json.charAt(s) == '"')) s++;
        int e = s;
        while (e < json.length() && json.charAt(e) != '"' && json.charAt(e) != ','
                && json.charAt(e) != '}') e++;
        return json.substring(s, e).trim();
    }

    static long amountOf(String json) {
        String a = field(json, "amount");
        try { return a == null ? 0L : Long.parseLong(a); } catch (Exception e) { return 0L; }
    }

    static String orderIdOf(String json) {
        String o = field(json, "order_id");
        return o == null ? "UNKNOWN" : o;
    }

    // Streams 앱을 띄우고 Ctrl+C 까지 대기합니다.
    // shutdown hook 으로 close() 를 걸어 두어야 컨슈머 그룹에서 깨끗이 빠집니다.
    // 안 그러면 13-11 의 reset 도구가 "still active" 로 거부합니다.
    static void start(Topology topology, Properties props) {
        System.out.println("[TOPOLOGY]\n" + topology.describe());
        KafkaStreams streams = new KafkaStreams(topology, props);
        CountDownLatch latch = new CountDownLatch(1);

        streams.setStateListener((now, old) ->
                System.out.printf("[STATE] %s → %s%n", old, now));
        streams.setUncaughtExceptionHandler(e -> {
            System.err.println("[FATAL] " + e);
            return StreamsUncaughtExceptionHandler.StreamThreadExceptionResponse.SHUTDOWN_CLIENT;
        });

        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            System.out.println("\n[SHUTDOWN] closing...");
            streams.close(Duration.ofSeconds(10));
            latch.countDown();
        }));

        streams.start();
        System.out.println("[START] application.id = " + props.get(StreamsConfig.APPLICATION_ID_CONFIG));
        System.out.println("[START] Ctrl+C 로 종료합니다.");
        try { latch.await(); } catch (InterruptedException ignored) { Thread.currentThread().interrupt(); }
    }

    // 타임스탬프를 직접 지정해서 보내는 프로듀서.
    // 콘솔 프로듀서로는 타임스탬프를 못 정하므로 grace 실습은 Java 로만 가능합니다.
    static void sendAt(String topic, String key, String value, long timestampMs) {
        Properties p = new Properties();
        p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
        p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        p.put(ProducerConfig.ACKS_CONFIG, "all");
        try (Producer<String, String> prod = new KafkaProducer<>(p)) {
            prod.send(new ProducerRecord<>(topic, null, timestampMs, key, value));
            prod.flush();
            System.out.printf("[PRODUCE] %s %s ts=%s%n", key, value, Instant.ofEpochMilli(timestampMs));
        }
    }

    interface Scenario { void run(String opt) throws Exception; }

    // ========================================================================
    // topology — 실행하지 않고 describe() 만 출력
    // ========================================================================
    static class TopologyOnly {
        static void run(String which) {
            StreamsBuilder b = new StreamsBuilder();
            switch (which.isEmpty() ? "stateless" : which) {
                case "selectkey" -> {
                    b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
                     .selectKey((k, v) -> orderIdOf(v))
                     .groupByKey()
                     .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                             .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
                     .toStream()
                     .to(T_OUT, Produced.with(Serdes.String(), Serdes.Long()));
                }
                case "count" -> {
                    b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
                     .groupByKey()
                     .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                             .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
                     .toStream()
                     .to(T_OUT, Produced.with(Serdes.String(), Serdes.Long()));
                }
                case "window" -> {
                    b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
                     .groupByKey()
                     .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
                     .aggregate(() -> 0L, (k, v, agg) -> agg + amountOf(v),
                             Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("amount-per-min")
                                     .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
                     .toStream()
                     .foreach((wk, sum) -> { });
                }
                case "join" -> {
                    KStream<String, String> orders   = b.stream(T_ORDERS);
                    KStream<String, String> payments = b.stream(T_PAYMENTS);
                    orders.selectKey((k, v) -> orderIdOf(v))
                          .join(payments,
                                  (o, p) -> "{\"order\":" + o + ",\"payment\":" + p + "}",
                                  JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30)),
                                  StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String()))
                          .to(T_OUT);
                }
                default -> {
                    b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
                     .filter((k, v) -> v.contains("CREATED"))
                     .mapValues(String::toUpperCase)
                     .to(T_OUT, Produced.with(Serdes.String(), Serdes.String()));
                }
            }
            System.out.println(b.build().describe());
            // 서브토폴로지가 2개 이상이면 그 경계마다 repartition 토픽이 끼어들고,
            // 데이터가 Kafka 를 한 번 왕복합니다. 즉 서브토폴로지 개수 = 네트워크 왕복 횟수.
        }
    }

    // ========================================================================
    // stateless — 상태를 안 쓰는 연산들
    // ========================================================================
    static class Stateless implements Scenario {
        public void run(String opt) {
            StreamsBuilder b = new StreamsBuilder();

            KStream<String, String> src = b.stream(T_ORDERS,
                    Consumed.with(Serdes.String(), Serdes.String()));

            // peek — 아무것도 안 바꾸고 들여다보기만 합니다. 디버깅용.
            KStream<String, String> seen =
                    src.peek((k, v) -> System.out.printf("[PEEK-in ] %s → %s%n", k, v));

            // filter — CANCELLED 는 버립니다.
            KStream<String, String> live = seen
                    .peek((k, v) -> {
                        if (v.contains("CANCELLED"))
                            System.out.printf("[filter  ] CANCELLED 제외됨: %s%n", orderIdOf(v));
                    })
                    .filter((k, v) -> !v.contains("CANCELLED"));

            // branch / split — 금액으로 두 갈래로 나눕니다.
            Map<String, KStream<String, String>> branches = live.split(Named.as("br-"))
                    .branch((k, v) -> amountOf(v) >= 10000, Branched.as("big"))
                    .defaultBranch(Branched.as("small"));

            KStream<String, String> big = branches.get("br-big")
                    .peek((k, v) -> System.out.printf("[branch  ] BIG   ← %s amount=%d%n", k, amountOf(v)));
            KStream<String, String> small = branches.get("br-small")
                    .peek((k, v) -> System.out.printf("[branch  ] SMALL ← %s amount=%d%n", k, amountOf(v)));

            // merge — 다시 하나로 합칩니다.
            // ★ mapValues 를 씁니다. map 이 아니라. map 은 키를 안 바꿔도
            //   "바뀌었을 수 있다"고 표시되어 다음 집계 앞에 repartition 을 유발합니다.
            big.merge(small)
               .mapValues(String::toUpperCase)
               .peek((k, v) -> System.out.printf("[merge   ] → %s%n", k))
               .to(T_OUT, Produced.with(Serdes.String(), Serdes.String()));

            // 이 토폴로지는 스테이트리스이므로 내부 토픽이 하나도 안 생깁니다.
            // kt --list 로 확인해 보세요.
            start(b.build(), baseProps("s13-stateless-app"));
        }
    }

    // ========================================================================
    // selectkey — 키를 바꾸면 repartition 토픽이 생깁니다
    // ========================================================================
    static class SelectKey implements Scenario {
        public void run(String opt) {
            StreamsBuilder b = new StreamsBuilder();

            b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
             // 키를 customer_id → order_id 로 바꿉니다.
             // ★ 이 순간 레코드는 여전히 "원래 파티션"에 있습니다. 키만 바뀌었습니다.
             //   같은 order_id 가 여러 파티션에 흩어져 있을 수 있으므로,
             //   집계 전에 Streams 가 repartition 토픽에 다시 써서 재분배합니다.
             .selectKey((k, v) -> orderIdOf(v))
             .groupByKey()
             .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                     .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
             .toStream()
             .peek((k, c) -> System.out.printf("[COUNT] %s = %d%n", k, c))
             .to(T_OUT, Produced.with(Serdes.String(), Serdes.Long()));

            // 실행 후 생기는 토픽:
            //   s13-selectkey-app-order-counts-repartition  (cleanup.policy=delete)
            //   s13-selectkey-app-order-counts-changelog    (cleanup.policy=compact)
            start(b.build(), baseProps("s13-selectkey-app"));
        }
    }

    // ========================================================================
    // count — 집계와 상태 저장소
    // ========================================================================
    static class Count implements Scenario {
        public void run(String opt) {
            StreamsBuilder b = new StreamsBuilder();

            // groupByKey() 는 키를 안 바꾸므로 repartition 토픽이 안 생깁니다.
            // groupBy((k,v) -> ...) 였다면 생깁니다. 이 차이를 기억하세요.
            b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
             .groupByKey()
             .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                     .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
             .toStream()
             // KTable 은 "변경로그"이므로 값이 바뀔 때마다 갱신을 내보냅니다.
             // C001 이 1, 2, 3 으로 세 번 나옵니다. 최종값 하나만 원하면 suppress 시나리오로.
             .peek((k, c) -> System.out.printf("[COUNT] %s = %d%n", k, c))
             .to(T_OUT, Produced.with(Serdes.String(), Serdes.Long()));

            // 상태는 두 곳에 있습니다:
            //   ① /tmp/kafka-streams/s13-count-app/0_0/rocksdb/order-counts/  (로컬, 빠른 읽기)
            //   ② s13-count-app-order-counts-changelog                        (Kafka, 진짜 원본)
            // 재기동 시 ① 이 없으면 ② 를 처음부터 읽어 복원합니다.
            start(b.build(), baseProps("s13-count-app"));
        }
    }

    // ========================================================================
    // window — 1분 텀블링 + grace period
    //   --inject-late 를 주면 grace 를 넘긴 레코드를 하나 넣어 "조용히 버려지는 것"을 재현
    // ========================================================================
    static class Window implements Scenario {
        public void run(String opt) throws Exception {
            boolean injectLate = "--inject-late".equals(opt);

            StreamsBuilder b = new StreamsBuilder();

            b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
             .groupByKey()
             // ofSizeAndGrace(크기, grace)
             //   grace = "윈도우가 끝난 뒤 얼마나 더 늦은 레코드를 받아 줄지"
             //   3.0 부터 grace 를 명시하도록 강제되었습니다.
             //   TimeWindows.of(size) 는 deprecated 이고 옛 기본 grace 는 24시간이었습니다.
             //   ofSizeWithNoGrace() 로 마이그레이션하면 grace 가 0 이 되어
             //   조금만 늦어도 전부 버려집니다. 3.x 마이그레이션의 대표 사고 지점입니다.
             .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
             .aggregate(() -> 0L,
                     (k, v, agg) -> agg + amountOf(v),
                     Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("amount-per-min")
                             .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
             .toStream()
             .foreach((wk, sum) -> System.out.printf("[WINDOW ] %s [%s ~ %s] sum=%d%n",
                     wk.key(), wk.window().startTime(), wk.window().endTime(), sum));

            Properties props = baseProps("s13-window-app");
            KafkaStreams streams = new KafkaStreams(b.build(), props);
            Runtime.getRuntime().addShutdownHook(new Thread(() -> streams.close(Duration.ofSeconds(10))));

            System.out.println("[TOPOLOGY]\n" + b.build().describe());
            streams.start();
            Thread.sleep(12_000);   // 리밸런싱 대기

            if (injectLate) {
                // 윈도우 경계는 에폭 기준 절대 시각입니다. 지금 시각을 1분 단위로 내림합니다.
                long base = (System.currentTimeMillis() / 60000L) * 60000L;

                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1001\",\"amount\":10000}", base + 10_000);
                Thread.sleep(3000);
                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1002\",\"amount\":20000}", base + 40_000);
                Thread.sleep(3000);

                // 스트림 시간을 윈도우 끝 + 20초로 밀어 올립니다. (grace 30초 안 → 아직 유효)
                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1003\",\"amount\":0}",     base + 80_000);
                Thread.sleep(2000);
                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1004\",\"amount\":7000}",  base + 50_000);
                Thread.sleep(3000);

                // 스트림 시간을 윈도우 끝 + 45초로 밀어 올립니다. (grace 30초 초과)
                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1005\",\"amount\":0}",     base + 105_000);
                Thread.sleep(2000);

                System.out.println("\n★ 이제 grace 를 넘긴 레코드를 넣습니다. 아무 출력도 없을 것입니다.");
                sendAt(T_ORDERS, "C001", "{\"order_id\":\"O-1006\",\"amount\":99000}", base + 55_000);
                Thread.sleep(5000);

                // dropped-records-total 지표를 읽습니다.
                // ★ 이 지표가 0 이 아니면 데이터가 조용히 버려지고 있다는 뜻입니다.
                //   예외도 안 나고 기본 로그 레벨에서는 아무것도 안 찍히므로,
                //   이 지표에 알람을 거는 것이 유일한 방어입니다.
                streams.metrics().forEach((name, metric) -> {
                    if ("dropped-records-total".equals(name.name())) {
                        System.out.printf("[METRIC ] %s (task=%s) = %s%n",
                                name.name(), name.tags().get("task-id"), metric.metricValue());
                    }
                });
                System.out.println("\n★ 99000 이 합계에 반영되지 않았습니다. 에러는 없었습니다.");
            } else {
                System.out.println("[HINT] 다른 터미널에서 kafka-console-producer.sh 로 s13_orders 에 넣어 보세요.");
                System.out.println("[HINT] grace 초과 재현은: Practice.java window --inject-late");
            }

            new CountDownLatch(1).await();
        }
    }

    // ========================================================================
    // suppress — 윈도우가 닫힌 뒤 최종 결과만
    // ========================================================================
    static class Suppress implements Scenario {
        public void run(String opt) throws Exception {
            StreamsBuilder b = new StreamsBuilder();

            b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
             .groupByKey()
             .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
             .aggregate(() -> 0L,
                     (k, v, agg) -> agg + amountOf(v),
                     Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("amount-suppressed")
                             .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
             // ⚠️ untilWindowCloses 는 "스트림 시간"이 윈도우끝+grace 를 넘어야 방출합니다.
             //    스트림 시간은 레코드가 들어와야 전진합니다. 트래픽이 멈추면
             //    마지막 윈도우는 영원히 방출되지 않습니다.
             //    "suppress 를 걸었더니 아무것도 안 나온다"의 열에 아홉이 이것입니다.
             //
             // ⚠️ BufferConfig.unbounded() 는 방출 전까지 모든 윈도우를 메모리에 들고 있습니다.
             //    키가 많으면 OOM 입니다. 운영에서는:
             //      Suppressed.BufferConfig.maxBytes(50_000_000L).emitEarlyWhenFull()
             .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()))
             .toStream()
             .foreach((wk, sum) -> System.out.printf("[FINAL  ] %s [%s ~ %s] sum=%d%n",
                     wk.key(), wk.window().startTime(), wk.window().endTime(), sum));

            System.out.println("[NOTE] 중간 결과는 안 나옵니다. 윈도우가 닫혀야 나옵니다.");
            System.out.println("[NOTE] 그리고 윈도우를 닫으려면 '다음 레코드'가 들어와야 합니다.");
            start(b.build(), baseProps("s13-suppress-app"));
        }
    }

    // ========================================================================
    // join — orders × payments 윈도우 조인
    // ========================================================================
    static class Join implements Scenario {
        public void run(String opt) {
            StreamsBuilder b = new StreamsBuilder();

            KStream<String, String> orders   = b.stream(T_ORDERS,
                    Consumed.with(Serdes.String(), Serdes.String()));
            KStream<String, String> payments = b.stream(T_PAYMENTS,
                    Consumed.with(Serdes.String(), Serdes.String()));

            // orders 의 키는 customer_id, payments 의 키는 order_id 입니다.
            // 조인하려면 키를 맞춰야 하므로 orders 쪽 키를 order_id 로 바꿉니다.
            // → 이 selectKey 때문에 repartition 토픽이 하나 생깁니다.
            orders.selectKey((k, v) -> orderIdOf(v))
                  .join(payments,
                          (o, p) -> "{\"order\":" + o + ",\"payment\":" + p + "}",
                          // KStream-KStream 조인은 윈도우가 필수입니다.
                          // "양쪽 레코드의 타임스탬프 차이가 5분 이내면 조인"이라는 뜻입니다.
                          JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30)),
                          StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String()))
                  .peek((k, v) -> System.out.printf("[JOIN] %s%n", k))
                  .to(T_OUT, Produced.with(Serdes.String(), Serdes.String()));

            // KStream-KStream 조인은 양쪽 모두 윈도우 스토어를 만듭니다.
            //   s13-join-app-KSTREAM-JOINTHIS-...-store-changelog
            //   s13-join-app-KSTREAM-JOINOTHER-...-store-changelog
            start(b.build(), baseProps("s13-join-app"));
        }
    }

    // ========================================================================
    // copartition — co-partitioning 위반 재현
    //   s13_left(3) × s13_right(6) 을 조인하면 TopologyException 이 납니다.
    //   --fix 를 주면 repartition 으로 파티션 수를 맞춘 버전이 돌아갑니다.
    // ========================================================================
    static class CoPartition implements Scenario {
        public void run(String opt) {
            boolean fix = "--fix".equals(opt);
            StreamsBuilder b = new StreamsBuilder();

            KStream<String, String> left = b.stream(T_LEFT,
                    Consumed.with(Serdes.String(), Serdes.String()));

            KStream<String, String> right;
            if (fix) {
                System.out.println("[FIX] s13_right(6) → repartition(3) 로 맞춥니다");
                // repartition() 은 지정한 파티션 수의 새 토픽을 만들고 거기로 다시 씁니다.
                // 토픽 이름: s13-copart-app-right-fixed-repartition
                right = b.<String, String>stream(T_RIGHT,
                                Consumed.with(Serdes.String(), Serdes.String()))
                         .repartition(Repartitioned.<String, String>as("right-fixed")
                                 .withNumberOfPartitions(3)
                                 .withKeySerde(Serdes.String())
                                 .withValueSerde(Serdes.String()));
            } else {
                System.out.println("[BROKEN] s13_left(3) × s13_right(6) 을 그대로 조인합니다");
                right = b.stream(T_RIGHT, Consumed.with(Serdes.String(), Serdes.String()));
            }

            left.join(right,
                      (l, r) -> "left=" + l + " right=" + r,
                      JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30)),
                      StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String()))
                .foreach((k, v) -> System.out.printf("[JOIN] %s  ← %s%n", k, v));

            Topology t = b.build();
            System.out.println("[TOPOLOGY]\n" + t.describe());

            // ★ 파티션 수 불일치는 build() 가 아니라 "파티션 할당 시점"에 터집니다.
            //   즉 start() 이후 첫 리밸런싱에서 이 예외가 나옵니다:
            //
            //   org.apache.kafka.streams.errors.TopologyException: Invalid topology:
            //     Following topics do not have the same number of partitions:
            //     [s13_left(3), s13_right(6)]
            //
            //   ⚠️ co-partitioning 의 세 조건 중 Streams 가 검증하는 것은 1번뿐입니다.
            //     1. 파티션 수가 같을 것              ← 검증됨 (위 예외)
            //     2. 키의 타입/직렬화가 같을 것        ← 검증 안 됨
            //     3. 파티셔너가 같을 것                ← 검증 안 됨. 조인이 그냥 안 됩니다.
            //   3번이 진짜 함정입니다. 예외도 로그도 없이 결과가 비어 있습니다.
            start(t, baseProps("s13-copart-app"));
        }
    }

    // ========================================================================
    // eos — exactly_once_v2
    // ========================================================================
    static class Eos implements Scenario {
        public void run(String opt) {
            StreamsBuilder b = new StreamsBuilder();

            b.stream(T_ORDERS, Consumed.with(Serdes.String(), Serdes.String()))
             .groupByKey()
             .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
                     .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()))
             .toStream()
             .peek((k, c) -> System.out.printf("[COUNT] %s = %d%n", k, c))
             .to(T_OUT, Produced.with(Serdes.String(), Serdes.Long()));

            Properties props = baseProps("s13-eos-app");
            // ★ 이 한 줄이 Step 07 에서 손으로 짠 트랜잭션 코드를 전부 대체합니다.
            props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);

            System.out.println("[CONFIG] processing.guarantee = exactly_once_v2");
            System.out.println("[CONFIG] → 내부적으로 강제되는 값:");
            System.out.println("           enable.idempotence          = true");
            System.out.println("           max.in.flight.requests...   = 5");
            System.out.println("           acks                        = all");
            System.out.println("           transactional.id            = s13-eos-app-<uuid>-<taskId>");
            System.out.println("           isolation.level             = read_committed");
            System.out.println("           commit.interval.ms          = 100   (기본 30000 에서 변경됨)");
            System.out.println();
            System.out.println("[NOTE] 보장 범위: 입력 오프셋 커밋 + 상태 갱신 + 출력 쓰기가 한 트랜잭션.");
            System.out.println("[NOTE] ⚠️ Kafka 안에서만 유효합니다. 처리 로직에서 외부 API 를 호출하면");
            System.out.println("       그것은 트랜잭션 밖이라 롤백돼도 취소되지 않습니다.");
            System.out.println("[NOTE] 다운스트림 컨슈머는 isolation.level=read_committed 여야 합니다.");
            System.out.println();

            start(b.build(), props);
        }
    }
}

practice.sh

본문 13-0 ~ 13-11 의 모든 명령을 절 번호 주석과 함께 담았습니다. Practice.java 를 컨테이너에 복사하고, 시나리오를 백그라운드로 띄우고, 데이터를 넣고, 결과를 확인하고, 마지막에 전부 정리합니다.

  • 맨 앞에서 docker cp Practice.java kafka-1:/tmp/ 를 한 번만 수행하고, 이후 모든 시나리오가 그 파일을 씁니다. 스크립트를 어느 디렉터리에서 실행해도 되도록 SCRIPT_DIRBASH_SOURCE 로 계산합니다.
  • run_app() 헬퍼가 시나리오를 백그라운드로 띄우고 PID 를 배열에 모읍니다. Streams 앱은 기동에 15~20초가 걸리므로(리밸런싱 + 내부 토픽 생성 + 상태 복원) 각 run_app 뒤에 sleep 20 이 붙어 있습니다. 이걸 줄이면 프로듀서가 먼저 데이터를 넣어 버려 앱이 못 받습니다.
  • [13-0] 에서 kt --list > /tmp/topics-before.txt시작 시점 토픽 목록을 기록합니다. [13-7]diff 가 이 파일을 씁니다. 이 before/after 대비가 함정 A 의 핵심 증거이므로 중간부터 실행하면 의미가 없습니다.
  • s13_right--partitions 6 으로 만드는 것이 의도적입니다. [13-9]copartition 시나리오가 TopologyException 을 내려면 파티션 수가 달라야 합니다.
  • [13-6]docker exec kafka-1 find /tmp/kafka-streams -maxdepth 4 -type dRocksDB 디렉터리를 직접 보여 줍니다. 0_0/0_1/0_2 가 태스크 ID(<서브토폴로지>_<파티션>)라는 것을 눈으로 확인하는 구간입니다.
  • [13-11] 정리가 세 단계입니다. ① 앱 종료(pkill -f Practice) ② kafka-streams-application-reset.sh 를 앱마다 실행 ③ rm -rf /tmp/kafka-streams. ③ 을 빼먹으면 다음 실행 때 옛 상태가 되살아나 카운트가 두 배가 됩니다. reset 도구가 로컬 상태를 안 지운다는 함정을 스크립트로 방어한 것입니다.
  • 컨슈머 그룹 삭제도 별도로 합니다. reset 도구는 오프셋만 0 으로 되돌리고 그룹은 남기기 때문입니다.
#!/usr/bin/env bash
set -euo pipefail
# ============================================================================
# Step 13 — Kafka Streams / practice.sh
#
# 실행법:
#   cd docs/reference/kafka/step-13-streams
#   bash practice.sh
#
#   한 단계씩 보려면:  bash -x practice.sh
#
# 이 스크립트는 Practice.java 를 kafka-1 컨테이너로 복사한 뒤
# 시나리오를 백그라운드로 띄우고, 데이터를 넣고, 결과를 확인하고, 정리합니다.
#
# ⚠️ Streams 앱은 기동에 15~20초가 걸립니다(리밸런싱 + 내부 토픽 생성 + 상태 복원).
#    각 run_app 뒤의 sleep 20 을 줄이면 프로듀서가 먼저 데이터를 넣어 버립니다.
# ============================================================================

BS=kafka-1:9092
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"

K()  { docker exec kafka-1 /opt/kafka/bin/"$@"; }
KT() { K kafka-topics.sh --bootstrap-server "$BS" "$@"; }
KCG(){ K kafka-consumer-groups.sh --bootstrap-server "$BS" "$@"; }
KGO(){ K kafka-get-offsets.sh --bootstrap-server "$BS" "$@"; }
KCONF(){ K kafka-configs.sh --bootstrap-server "$BS" "$@"; }
hr() { echo; echo "=============================================================="; echo "  $*"; echo "=============================================================="; }

APP_PIDS=()

run_app() {   # run_app <scenario> [opt]
  echo "--- 앱 기동: $* ---"
  docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java $*" &
  APP_PIDS+=("$!")
  sleep 20
}

stop_apps() {
  docker exec kafka-1 pkill -f 'Practice' >/dev/null 2>&1 || true
  for p in "${APP_PIDS[@]:-}"; do kill "$p" >/dev/null 2>&1 || true; done
  APP_PIDS=()
  sleep 6
}

produce() {   # produce <topic>  ← stdin 으로 key:value 줄
  docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
    --bootstrap-server "$BS" --topic "$1" \
    --property parse.key=true --property key.separator=:
}

consume() {   # consume <topic> [extra args...]
  local t="$1"; shift
  K kafka-console-consumer.sh --bootstrap-server "$BS" --topic "$t" \
    --from-beginning --property print.key=true --timeout-ms 5000 "$@" || true
}

# ----------------------------------------------------------------------------
# [13-0] 실습 준비
# ----------------------------------------------------------------------------
hr "[13-0] 실습 준비"

docker cp "$SCRIPT_DIR/Practice.java" kafka-1:/tmp/
echo "--- kafka-streams 가 배포판에 들어 있는지 확인 ---"
docker exec kafka-1 ls /opt/kafka/libs/ | grep -E 'kafka-streams|rocksdb'

echo "--- 사용법 출력 ---"
docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java" || true

echo "--- 실습 토픽 생성 ---"
KT --create --topic s13_orders   --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_payments --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_out      --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_left     --partitions 3 --replication-factor 3 --if-not-exists
# ★ s13_right 만 6파티션입니다. [13-9] 의 co-partitioning 함정을 위한 의도적 설정.
KT --create --topic s13_right    --partitions 6 --replication-factor 3 --if-not-exists

# ★ 시작 시점 토픽 목록을 기록합니다. [13-7] 의 diff 가 이 파일을 씁니다.
#   중간부터 실행하면 before/after 대비가 무의미해집니다.
KT --list | docker exec -i kafka-1 sh -c 'cat > /tmp/topics-before.txt'
echo "--- before ---"
docker exec kafka-1 cat /tmp/topics-before.txt

# ----------------------------------------------------------------------------
# [13-2] Topology — describe() 읽기
# ----------------------------------------------------------------------------
hr "[13-2] Topology"
docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java topology"
# 서브토폴로지가 1개 = 데이터가 네트워크를 안 거칩니다.
# 서브토폴로지 개수 = Kafka 네트워크 왕복 횟수입니다.

# ----------------------------------------------------------------------------
# [13-4] 스테이트리스 연산
#   filter / mapValues / branch / merge / peek. 내부 토픽이 하나도 안 생깁니다.
# ----------------------------------------------------------------------------
hr "[13-4] 스테이트리스"
run_app stateless

produce s13_orders <<'EOF'
C001:{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
C002:{"order_id":"O-1002","customer_id":"C002","amount":500,"status":"CREATED"}
C001:{"order_id":"O-1003","customer_id":"C001","amount":88000,"status":"CANCELLED"}
EOF
sleep 6
echo "--- s13_out (CANCELLED 는 빠지고 2건) ---"
consume s13_out

echo "--- 토픽 목록: 아직 안 늘었습니다 ---"
KT --list
stop_apps

# ----------------------------------------------------------------------------
# [13-5] selectKey → repartition 토픽이 생깁니다
# ----------------------------------------------------------------------------
hr "[13-5] selectKey 와 repartition"

echo "--- 먼저 토폴로지만 봅니다. 서브토폴로지가 둘로 쪼개집니다. ---"
docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java topology selectkey"

run_app selectkey
produce s13_orders <<'EOF'
C001:{"order_id":"O-2001","customer_id":"C001","amount":39000}
C002:{"order_id":"O-2002","customer_id":"C002","amount":12500}
EOF
sleep 6

echo "--- 토픽 목록: 두 개가 생겼습니다 ---"
KT --list | grep 's13-selectkey-app' || true

echo "--- repartition 토픽 (cleanup.policy=delete, RF 는 3 — baseProps 에서 지정했으므로) ---"
KT --describe --topic s13-selectkey-app-order-counts-repartition
# ⚠️ REPLICATION_FACTOR_CONFIG 을 안 줬다면 여기 ReplicationFactor 가 1 입니다.
#    Streams 기본값이 1 이라서입니다. 브로커 기본값(3)을 따르지 않습니다.
stop_apps

# ----------------------------------------------------------------------------
# [13-6] 집계와 상태 저장소
# ----------------------------------------------------------------------------
hr "[13-6] 집계와 상태 저장소"
run_app count

produce s13_orders <<'EOF'
C001:{"order_id":"O-3001","amount":39000}
C002:{"order_id":"O-3002","amount":12500}
C001:{"order_id":"O-3003","amount":88000}
C001:{"order_id":"O-3004","amount":5000}
EOF
sleep 6

echo "--- s13_out: C001 이 1,2,3 으로 세 번 나옵니다 (KTable = 변경로그) ---"
K kafka-console-consumer.sh --bootstrap-server "$BS" --topic s13_out --from-beginning \
  --property print.key=true \
  --value-deserializer org.apache.kafka.common.serialization.LongDeserializer \
  --timeout-ms 5000 || true

echo "--- 로컬 RocksDB 디렉터리 (0_0/0_1/0_2 는 태스크 ID = <서브토폴로지>_<파티션>) ---"
docker exec kafka-1 find /tmp/kafka-streams -maxdepth 4 -type d | head -12

echo "--- changelog 토픽: cleanup.policy=compact 입니다 (repartition 과 다름) ---"
KT --describe --topic s13-count-app-order-counts-changelog

echo "--- changelog 내용 ---"
K kafka-console-consumer.sh --bootstrap-server "$BS" \
  --topic s13-count-app-order-counts-changelog --from-beginning \
  --property print.key=true \
  --value-deserializer org.apache.kafka.common.serialization.LongDeserializer \
  --timeout-ms 5000 || true

echo "--- application.id 가 곧 컨슈머 그룹 ID 입니다 ---"
KCG --list
KCG --describe --group s13-count-app
stop_apps

# ----------------------------------------------------------------------------
# [13-8] 윈도우 + grace period (함정 B)
# ----------------------------------------------------------------------------
hr "[13-8] 윈도우와 grace"
echo "--- grace 초과 레코드를 주입합니다. Java 로만 가능합니다(타임스탬프 지정). ---"
docker exec kafka-1 sh -c "cd /tmp && timeout 90 java -cp '/opt/kafka/libs/*' Practice.java window --inject-late" || true
# ★ 마지막 99000 이 합계에 안 들어갑니다. 예외도 로그도 없습니다.
#   dropped-records-total 지표가 1 인 것만이 증거입니다.
stop_apps

hr "[13-8b] suppress — 최종 결과만"
echo "--- suppress 는 '다음 레코드'가 와야 방출합니다. 저트래픽에서 안 나오는 이유입니다. ---"
run_app suppress
produce s13_orders <<'EOF'
C001:{"order_id":"O-4001","amount":10000}
C002:{"order_id":"O-4002","amount":5000}
EOF
sleep 10
echo "(윈도우가 안 닫혀서 아직 아무것도 안 나옵니다 — 이게 정상입니다)"
stop_apps

# ----------------------------------------------------------------------------
# [13-9] 조인
# ----------------------------------------------------------------------------
hr "[13-9] 조인"
run_app join

produce s13_orders <<'EOF'
C001:{"order_id":"O-5001","customer_id":"C001","amount":39000,"status":"CREATED"}
EOF
sleep 3
produce s13_payments <<'EOF'
O-5001:{"order_id":"O-5001","method":"CARD","amount":39000,"result":"APPROVED"}
EOF
sleep 6
echo "--- 조인 결과 ---"
consume s13_out
stop_apps

hr "[13-9b] 함정 C — co-partitioning 위반"
echo "--- s13_left(3) × s13_right(6) → TopologyException ---"
docker exec kafka-1 sh -c "cd /tmp && timeout 45 java -cp '/opt/kafka/libs/*' Practice.java copartition" 2>&1 | tail -12 || true
stop_apps

echo "--- --fix: repartition 으로 파티션 수를 맞춥니다 ---"
docker exec kafka-1 sh -c "cd /tmp && timeout 45 java -cp '/opt/kafka/libs/*' Practice.java copartition --fix" 2>&1 | head -30 || true
stop_apps

# ----------------------------------------------------------------------------
# [13-10] EOS
# ----------------------------------------------------------------------------
hr "[13-10] exactly_once_v2"
docker exec kafka-1 sh -c "cd /tmp && timeout 45 java -cp '/opt/kafka/libs/*' Practice.java eos" 2>&1 | head -25 || true
stop_apps

# ----------------------------------------------------------------------------
# [13-7] 함정 A — 내부 토픽이 자동으로 생깁니다 (before/after 비교)
# ----------------------------------------------------------------------------
hr "[13-7] 함정 A — 자동 생성된 내부 토픽"
KT --list | docker exec -i kafka-1 sh -c 'cat > /tmp/topics-after.txt'
echo "--- diff (좌: before / 우: after) ---"
docker exec kafka-1 sh -c 'diff /tmp/topics-before.txt /tmp/topics-after.txt' || true

echo "--- 그런데 자동 토픽 생성은 꺼져 있습니다 ---"
KCONF --describe --entity-type brokers --entity-name 1 | grep auto.create || true
# ★ Streams 는 프로듀서/컨슈머의 자동 생성 경로를 안 씁니다.
#   InternalTopicManager 가 AdminClient 로 CreateTopics API 를 직접 호출합니다.
#   사용자가 손으로 kt --create 하는 것과 완전히 같은 경로라 auto.create 로 못 막습니다.
#   Connect(Step 12)의 내부 토픽도 같은 방식입니다.

# ----------------------------------------------------------------------------
# [13-11] 정리 — 세 단계입니다
#   ① 앱 종료  ② reset 도구(Kafka 쪽)  ③ 로컬 상태 디렉터리(★ 빼먹기 쉬움)
# ----------------------------------------------------------------------------
hr "[13-11] 정리"

echo "--- ① 앱 전부 종료 ---"
stop_apps
KCG --list

echo "--- ② kafka-streams-application-reset.sh ---"
for app in s13-stateless-app s13-selectkey-app s13-count-app s13-window-app \
           s13-suppress-app s13-join-app s13-copart-app s13-eos-app; do
  echo "  reset: $app"
  K kafka-streams-application-reset.sh --bootstrap-server "$BS" \
     --application-id "$app" --input-topics s13_orders 2>&1 | tail -3 || true
done

echo "--- 컨슈머 그룹 삭제 (reset 도구는 오프셋만 되돌리고 그룹은 남깁니다) ---"
for g in $(KCG --list | grep '^s13-' || true); do KCG --delete --group "$g" || true; done
KCG --list

echo "--- ③ 로컬 상태 디렉터리 (★ reset 도구가 안 지웁니다) ---"
docker exec kafka-1 du -sh /tmp/kafka-streams 2>/dev/null || echo "없음"
docker exec kafka-1 rm -rf /tmp/kafka-streams
docker exec kafka-1 ls /tmp/kafka-streams 2>&1 || echo "삭제 완료"
# ⚠️ ③ 을 빼먹으면 다음 실행 때 옛 RocksDB 상태가 되살아나
#    "오프셋은 0 인데 카운트는 이어지는" 상태가 됩니다. 에러는 없습니다.

echo "--- 실습 토픽 삭제 ---"
for t in $(KT --list | grep -E '^s13' || true); do KT --delete --topic "$t" || true; done
KT --list | grep -E '^s13' || echo "s13 토픽 없음"

docker exec kafka-1 rm -f /tmp/Practice.java /tmp/topics-before.txt /tmp/topics-after.txt

hr "practice.sh 완료"

exercise.sh

6문제의 문제지입니다. 각 문제는 # 여기에 작성: 자리를 비워 두었고, 관찰에 필요한 앱은 문제지 쪽에서 미리 띄워 줍니다.

  • 문제 1·2·3 은 내부 토픽을 관찰·판별하는 문제, 문제 4·5·6 은 직접 조치하는 문제입니다.
  • 문제 1 은 statelessselectkey 두 시나리오를 차례로 띄우고 그 사이사이에 kt --list 를 찍습니다. 여러분이 할 일은 세 스냅숏을 비교해 어느 연산 뒤에 토픽이 생겼는지 짚고, 생긴 토픽 이름의 각 부분이 무엇을 뜻하는지 쓰는 것입니다.
  • 문제 2 는 토픽 이름 5개를 주고 소유 앱과 종류(repartition/changelog)를 맞히는 순수 지필 문제입니다. 그중 하나는 KSTREAM-AGGREGATE-STATE-STORE-0000000003 처럼 자동 생성 이름이라, "이런 이름이 나오면 Materialized.as() 를 안 준 것"이라는 결론까지 끌어내는 것이 목표입니다.
  • 문제 4 는 함정이 있습니다. kafka-streams-application-reset.sh 만 돌리고 앱을 다시 띄우면 카운트가 이어집니다. 로컬 state.dir 을 안 지웠기 때문입니다. 문제지가 그 상태까지 만들어 두고 "왜 카운트가 1부터 시작 안 합니까?"를 묻습니다.
  • 문제 5 는 s13ex_a(2파티션)와 s13ex_b(4파티션)를 만들어 둡니다. 예외를 재현한 뒤 repartition(Repartitioned.numberOfPartitions(2)) 로 고치는 것이 정답이며, Practice.java copartition --fix 를 참고 구현으로 볼 수 있습니다.
  • 문제 6 은 Practice.java window --inject-late 를 실행해 둔 상태에서 시작합니다. 출력의 합계가 실제 투입 금액보다 작다는 것을 확인하고, 어느 JMX 지표를 봐야 하는지(dropped-records-total) 와 어느 설정을 고쳐야 하는지(grace 확대)를 답하면 됩니다.
  • 파일 끝의 정리 블록은 앱 종료 → reset → 그룹 삭제 → rm -rf /tmp/kafka-streamss13ex_ 토픽 삭제 순서이며, 문제를 안 풀고 정리만 돌려도 에러가 안 나도록 전부 || true 가 붙어 있습니다.
#!/usr/bin/env bash
set -euo pipefail
# ============================================================================
# Step 13 — Kafka Streams / exercise.sh   (문제 6개)
#
# 실행법:
#   cd docs/reference/kafka/step-13-streams
#   bash exercise.sh
#
# 각 문제의 "# 여기에 작성:" 아래를 직접 채우세요.
# 정답은 solution.sh 에 있습니다. 먼저 풀어 본 뒤에 여세요.
#
# 문제 1·2·3 = 내부 토픽을 관찰·판별하는 문제
# 문제 4·5·6 = 직접 조치하는 문제
# ============================================================================

BS=kafka-1:9092
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"

K()  { docker exec kafka-1 /opt/kafka/bin/"$@"; }
KT() { K kafka-topics.sh --bootstrap-server "$BS" "$@"; }
KCG(){ K kafka-consumer-groups.sh --bootstrap-server "$BS" "$@"; }
hr() { echo; echo "------------------------------------------------------------"; echo "  $*"; echo "------------------------------------------------------------"; }

run_app() {
  docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java $*" &
  sleep 20
}
stop_apps() { docker exec kafka-1 pkill -f 'Practice' >/dev/null 2>&1 || true; sleep 6; }
produce() {
  docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
    --bootstrap-server "$BS" --topic "$1" \
    --property parse.key=true --property key.separator=:
}

docker cp "$SCRIPT_DIR/Practice.java" kafka-1:/tmp/
KT --create --topic s13_orders   --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_payments --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_out      --partitions 3 --replication-factor 3 --if-not-exists


# ============================================================================
# 문제 1. 어느 연산 뒤에 내부 토픽이 생깁니까?
#
#   아래 준비 블록이 ① 아무것도 안 띄운 상태 ② stateless 앱 ③ selectkey 앱
#   세 시점의 토픽 목록을 찍어 줍니다.
#
#   (a) 어느 시점에 토픽이 늘었습니까? 늘어난 토픽 이름을 쓰세요.
#   (b) 그 이름을 세 조각으로 나눠 각각 무엇을 뜻하는지 쓰세요.
#   (c) stateless 앱은 왜 아무 토픽도 안 만듭니까?
# ============================================================================
hr "문제 1 — 준비"
echo "=== ① 시작 시점 ==="
KT --list | grep -v '^__' || true

echo "=== ② stateless 앱 기동 후 ==="
run_app stateless
produce s13_orders <<'EOF'
C001:{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
EOF
sleep 5
KT --list | grep -v '^__' || true
stop_apps

echo "=== ③ selectkey 앱 기동 후 ==="
run_app selectkey
produce s13_orders <<'EOF'
C001:{"order_id":"O-1002","customer_id":"C001","amount":11000}
EOF
sleep 6
KT --list | grep -v '^__' || true
stop_apps

hr "문제 1 — 답"
# 여기에 작성: (a) 늘어난 토픽 이름


# 여기에 작성: (b) 이름의 세 조각과 각각의 의미


# 여기에 작성: (c) stateless 앱이 토픽을 안 만드는 이유



# ============================================================================
# 문제 2. 토픽 이름만 보고 소유 앱과 종류 맞히기
#
#   운영 클러스터에서 아래 다섯 개의 토픽을 발견했습니다.
#   각각 (소유 앱 / repartition·changelog / 어떤 연산 때문에 생겼는지) 를 쓰세요.
#
#     1) payment-agg-app-daily-total-changelog
#     2) payment-agg-app-daily-total-repartition
#     3) order-join-app-KSTREAM-JOINTHIS-0000000012-store-changelog
#     4) fraud-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
#     5) fraud-app-KSTREAM-KEY-SELECT-0000000002-repartition
#
#   (보너스) 4)·5) 같은 이름이 나오는 앱은 코드에 무엇이 빠져 있습니까?
#            그것이 배포할 때 왜 위험합니까?
# ============================================================================
hr "문제 2"
# 여기에 작성: 다섯 개 각각의 소유 앱 / 종류 / 원인 연산


# 여기에 작성: (보너스) 4)·5) 의 앱에 빠진 것과 배포 시 위험



# ============================================================================
# 문제 3. changelog 토픽의 cleanup.policy
#
#   (a) count 앱을 띄우고 changelog 토픽의 cleanup.policy 를 확인하는 명령을 쓰세요.
#   (b) 그 값이 delete 였다면 무슨 일이 생깁니까?
#   (c) repartition 토픽은 왜 delete 여도 괜찮습니까?
# ============================================================================
hr "문제 3 — 준비"
run_app count
produce s13_orders <<'EOF'
C001:{"order_id":"O-2001","amount":1000}
C002:{"order_id":"O-2002","amount":2000}
EOF
sleep 6
stop_apps

hr "문제 3 — 답"
# 여기에 작성: (a) 확인 명령


# 여기에 작성: (b) delete 였다면


# 여기에 작성: (c) repartition 이 delete 여도 되는 이유



# ============================================================================
# 문제 4. 앱을 완전히 초기화하기
#
#   아래 준비 블록이 count 앱으로 4건을 처리해 카운트를 만들어 둡니다.
#   여러분은 이 앱을 "완전히 처음 상태"로 되돌린 뒤 다시 띄워서
#   카운트가 1부터 시작하게 만들어야 합니다.
#
#   ⚠️ kafka-streams-application-reset.sh 만 돌리고 다시 띄우면
#      카운트가 이어집니다. 왜 그렇습니까? 무엇을 더 해야 합니까?
# ============================================================================
hr "문제 4 — 준비"
run_app count
produce s13_orders <<'EOF'
C001:{"order_id":"O-3001","amount":1000}
C001:{"order_id":"O-3002","amount":2000}
C001:{"order_id":"O-3003","amount":3000}
C001:{"order_id":"O-3004","amount":4000}
EOF
sleep 6
echo "지금 C001 의 카운트가 4 이상이어야 합니다."
stop_apps

hr "문제 4 — 답"
# 여기에 작성: 앱을 완전히 초기화하는 명령 (세 줄)


# 여기에 작성: 재기동해서 카운트가 1부터 시작하는지 확인


# 여기에 작성: reset 도구만으로는 왜 부족합니까? (주석)



# ============================================================================
# 문제 5. co-partitioning 위반 재현과 수정
#
#   s13ex_a(2파티션)과 s13ex_b(4파티션)을 만들어 두었습니다.
#   (a) 이 둘을 조인하면 어떤 예외가 납니까? (Practice.java copartition 을 참고)
#   (b) 예외 메시지를 그대로 쓰세요.
#   (c) 어떻게 고칩니까? 코드 한 줄로 답하세요.
#   (d) co-partitioning 의 조건 세 가지 중 Streams 가 검증해 주는 것은 몇 개입니까?
#       검증 안 되는 것이 왜 더 위험합니까?
# ============================================================================
hr "문제 5 — 준비"
KT --create --topic s13ex_a --partitions 2 --replication-factor 3 --if-not-exists
KT --create --topic s13ex_b --partitions 4 --replication-factor 3 --if-not-exists
KT --describe --topic s13ex_a | head -1
KT --describe --topic s13ex_b | head -1

hr "문제 5 — 답"
# 여기에 작성: (a)(b) 예외를 재현하고 메시지 확인
#   힌트: Practice.java 의 copartition 시나리오가 s13_left(3) × s13_right(6) 을 씁니다.
#         같은 원리입니다.


# 여기에 작성: (c) 고치는 코드 한 줄


# 여기에 작성: (d) 세 조건 중 검증되는 것 / 안 되는 것의 위험



# ============================================================================
# 문제 6. 윈도우 집계 결과가 실제보다 작습니다
#
#   아래 준비 블록이 window --inject-late 를 실행합니다.
#   총 투입 금액은 10000 + 20000 + 7000 + 99000 = 136000 인데
#   최종 합계는 37000 으로 나옵니다.
#
#   (a) 없어진 99000 은 어디로 갔습니까?
#   (b) 어느 지표를 보면 이것을 감지할 수 있습니까? 전체 지표 이름을 쓰세요.
#   (c) 어떻게 고칩니까? 코드에서 무엇을 바꿉니까?
#   (d) 그 값을 얼마로 잡아야 합니까? 판단 기준은?
# ============================================================================
hr "문제 6 — 준비"
docker exec kafka-1 sh -c "cd /tmp && timeout 90 java -cp '/opt/kafka/libs/*' Practice.java window --inject-late" 2>&1 | tail -20 || true
stop_apps

hr "문제 6 — 답"
# 여기에 작성: (a) 99000 의 행방


# 여기에 작성: (b) 감지할 지표 이름


# 여기에 작성: (c) 고치는 방법


# 여기에 작성: (d) 값을 정하는 기준



# ============================================================================
# 정리 — 문제를 안 풀고 정리만 돌려도 에러 없이 끝납니다.
# ============================================================================
hr "정리"
stop_apps
for app in s13-stateless-app s13-selectkey-app s13-count-app s13-window-app \
           s13-suppress-app s13-join-app s13-copart-app s13-eos-app s13ex-app; do
  K kafka-streams-application-reset.sh --bootstrap-server "$BS" \
     --application-id "$app" --input-topics s13_orders >/dev/null 2>&1 || true
done
for g in $(KCG --list | grep -E '^s13' || true); do KCG --delete --group "$g" >/dev/null 2>&1 || true; done
docker exec kafka-1 rm -rf /tmp/kafka-streams
for t in $(KT --list | grep -E '^s13' || true); do KT --delete --topic "$t" >/dev/null 2>&1 || true; done
docker exec kafka-1 rm -f /tmp/Practice.java
echo "정리 완료"

solution.sh

6문제의 정답 명령과 긴 해설 주석입니다. 문제를 풀어 본 뒤에 여세요.

  • 정답 1diff 결과가 <app-id>-order-counts-repartition<app-id>-order-counts-changelog 두 개라는 것이고, 해설에서 이름을 세 조각(앱ID / 스토어 또는 노드 이름 / 종류)으로 분해합니다. stateless 시나리오에서는 아무 토픽도 안 생긴다는 대비가 답의 절반입니다.
  • 정답 2 는 표로 답합니다. -repartition 은 키가 바뀐 뒤 다시 뿌리려고, -changelog 는 상태 저장소를 백업하려고 생깁니다. 자동 생성 이름(KSTREAM-...-0000000003)의 숫자가 DSL 호출 순서이므로, 코드에 filter 하나만 추가해도 번호가 밀려 토픽 이름이 바뀌고 기존 상태가 고아가 된다는 것이 이 문제의 진짜 교훈입니다.
  • 정답 3cleanup.policy=compact 이고, delete 였다면 retention 이 지난 뒤 앱을 재기동할 때 상태 일부가 사라진 채 복원되어 카운트가 조용히 작아진다는 설명입니다. repartition 토픽은 반대로 delete 가 맞는 이유(Streams 가 소비 직후 purge 하므로)도 나란히 씁니다.
  • 정답 4 의 핵심은 세 줄입니다. pkillkafka-streams-application-reset.shrm -rf /tmp/kafka-streams/<app-id>. 세 번째가 없으면 오프셋만 0 이 되고 로컬 RocksDB 는 남아 옛 상태에 새 데이터가 얹힙니다. 앱이 살아 있을 때 reset 을 돌리면 나오는 Consumer group ... is still active 예외 전문도 주석에 넣었습니다.
  • 정답 5 는 예외 메시지(Following topics do not have the same number of partitions: [s13ex_a(2), s13ex_b(4)])를 먼저 보여 주고, repartition(Repartitioned.numberOfPartitions(2)) 로 고칩니다. 그리고 co-partitioning 의 세 조건 중 Streams 가 검증하는 것은 1번(파티션 수)뿐이며, 3번(파티셔너 불일치)은 예외 없이 조인이 그냥 안 되는 침묵의 실패라는 점을 길게 설명합니다.
  • 정답 6kafka.streams:type=stream-task-metricsdropped-records-total 을 JmxTool 로 읽는 명령과, grace 를 Duration.ofSeconds(30)Duration.ofMinutes(5) 로 늘리는 코드 변경입니다. 해설에서 grace 를 "관측된 최대 지연의 1.5~2배"로 잡는 기준과, 3.0 에서 TimeWindows.of() 가 deprecated 되며 마이그레이션 시 grace 가 0 이 되어 버리는 사고를 경고합니다.
#!/usr/bin/env bash
set -euo pipefail
# ============================================================================
# Step 13 — Kafka Streams / solution.sh   (정답 + 해설)
#
# 실행법:
#   cd docs/reference/kafka/step-13-streams
#   bash solution.sh
#
# exercise.sh 를 먼저 풀어 본 뒤에 여세요.
# ============================================================================

BS=kafka-1:9092
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"

K()  { docker exec kafka-1 /opt/kafka/bin/"$@"; }
KT() { K kafka-topics.sh --bootstrap-server "$BS" "$@"; }
KCG(){ K kafka-consumer-groups.sh --bootstrap-server "$BS" "$@"; }
hr() { echo; echo "============================================================"; echo "  $*"; echo "============================================================"; }

run_app() {
  docker exec kafka-1 sh -c "cd /tmp && java -cp '/opt/kafka/libs/*' Practice.java $*" &
  sleep 20
}
stop_apps() { docker exec kafka-1 pkill -f 'Practice' >/dev/null 2>&1 || true; sleep 6; }
produce() {
  docker exec -i kafka-1 /opt/kafka/bin/kafka-console-producer.sh \
    --bootstrap-server "$BS" --topic "$1" \
    --property parse.key=true --property key.separator=:
}

docker cp "$SCRIPT_DIR/Practice.java" kafka-1:/tmp/
KT --create --topic s13_orders   --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_payments --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_out      --partitions 3 --replication-factor 3 --if-not-exists


# ============================================================================
# 정답 1 — 어느 연산 뒤에 내부 토픽이 생기는가
# ============================================================================
hr "정답 1"

KT --list | grep -v '^__' | sort > /tmp/sol-before.txt 2>/dev/null || true
docker exec kafka-1 sh -c "true"   # noop

echo "--- ② stateless 앱 ---"
run_app stateless
produce s13_orders <<'EOF'
C001:{"order_id":"O-1001","customer_id":"C001","amount":39000,"status":"CREATED"}
EOF
sleep 5
echo "s13-stateless-app 로 시작하는 토픽:"
KT --list | grep 's13-stateless-app' || echo "  (없음 — 이것이 정답의 절반입니다)"
stop_apps

echo "--- ③ selectkey 앱 ---"
run_app selectkey
produce s13_orders <<'EOF'
C001:{"order_id":"O-1002","customer_id":"C001","amount":11000}
EOF
sleep 6
echo "s13-selectkey-app 로 시작하는 토픽:"
KT --list | grep 's13-selectkey-app' || true
stop_apps

# 해설:
#   (a) 늘어난 토픽 두 개:
#         s13-selectkey-app-order-counts-repartition
#         s13-selectkey-app-order-counts-changelog
#
#   (b) 이름은 세 조각입니다:
#         s13-selectkey-app  |  order-counts  |  repartition
#         ─────────────────     ────────────     ───────────
#         application.id        상태 저장소 이름   토픽 종류
#                               (Materialized.as())
#
#       repartition — 키가 바뀐 뒤 같은 키를 같은 파티션으로 다시 모으려고 생깁니다.
#                     selectKey/map/groupBy 뒤에 집계나 조인이 오면 생깁니다.
#                     cleanup.policy=delete. Streams 가 소비 직후 스스로 purge 합니다.
#       changelog   — 상태 저장소(RocksDB)의 모든 변경을 Kafka 에 백업합니다.
#                     count/reduce/aggregate/table() 등 상태를 쓰면 생깁니다.
#                     cleanup.policy=compact. 이것이 상태의 "진짜 원본"입니다.
#
#   (c) stateless 앱은 왜 아무것도 안 만듭니까:
#       filter / mapValues / branch / merge / peek 은 전부 "레코드 하나만 보고" 처리합니다.
#       - 상태를 안 쓰므로 changelog 가 필요 없고
#       - 키를 안 바꾸므로(mapValues 를 썼습니다) 재분배가 필요 없어 repartition 도 없습니다
#       그래서 서브토폴로지도 1개이고, 데이터가 Kafka 를 왕복하지 않습니다.
#
#       ★ 만약 stateless 앱에서 mapValues 대신 map 을 썼다면?
#         map 은 "키가 바뀌었을 수 있다"고 표시됩니다. 실제로 키를 안 바꿔도 그렇습니다.
#         그 뒤에 집계가 오면 repartition 이 생깁니다. 처리량이 대략 절반이 됩니다.
#         → 항상 mapValues / flatMapValues / transformValues 를 먼저 고려하세요.


# ============================================================================
# 정답 2 — 토픽 이름으로 소유 앱과 종류 판별하기
# ============================================================================
hr "정답 2"
cat <<'ANSWER'
  1) payment-agg-app-daily-total-changelog
     앱: payment-agg-app / 종류: changelog
     원인: "daily-total" 이라는 이름의 상태 저장소를 쓰는 집계
           (count / reduce / aggregate 중 하나. Materialized.as("daily-total"))

  2) payment-agg-app-daily-total-repartition
     앱: payment-agg-app / 종류: repartition
     원인: 위 집계 앞에 키를 바꾸는 연산이 있음 (selectKey / map / groupBy)

  3) order-join-app-KSTREAM-JOINTHIS-0000000012-store-changelog
     앱: order-join-app / 종류: changelog
     원인: KStream-KStream 조인. 조인은 양쪽 각각에 윈도우 스토어를 만들므로
           JOINTHIS(왼쪽) 와 JOINOTHER(오른쪽) 가 쌍으로 생깁니다.
           이 토픽이 있으면 반드시 ...JOINOTHER-...-store-changelog 도 있습니다.

  4) fraud-app-KSTREAM-AGGREGATE-STATE-STORE-0000000003-changelog
     앱: fraud-app / 종류: changelog
     원인: 집계. 단 상태 저장소에 "이름을 안 줬습니다."

  5) fraud-app-KSTREAM-KEY-SELECT-0000000002-repartition
     앱: fraud-app / 종류: repartition
     원인: selectKey. 역시 이름 없음.
ANSWER

# 해설 (보너스):
#   4)·5) 의 앱에는 Materialized.as("이름") 이 빠져 있습니다.
#   이름을 안 주면 Streams 가 노드 이름을 씁니다:
#       KSTREAM-AGGREGATE-STATE-STORE-0000000003
#                                     ~~~~~~~~~~
#   ★ 이 숫자는 "DSL 호출 순서"입니다.
#
#   왜 위험합니까:
#     코드에 filter 를 하나 추가하면 그 뒤의 모든 노드 번호가 밀립니다.
#       0000000003 → 0000000004
#     토픽 이름이 바뀝니다. 배포하면
#       - 새 이름의 changelog 토픽이 새로 만들어지고
#       - 기존 상태가 담긴 옛 토픽은 아무도 안 읽는 고아가 되고
#       - 집계가 0 부터 다시 시작합니다
#     "필터 한 줄 추가했는데 집계가 리셋됐다"는 사고가 이렇게 납니다.
#     에러는 없습니다. 로그도 안 나옵니다. 숫자만 이상해집니다.
#
#   해결: 상태 저장소에는 예외 없이 이름을 주세요.
#       .count(Materialized.as("order-counts"))
#       .join(..., StreamJoined.as("order-payment-join"))
#       .repartition(Repartitioned.as("by-order-id"))
#   그러면 코드를 아무리 고쳐도 토픽 이름이 안 바뀝니다.


# ============================================================================
# 정답 3 — changelog 의 cleanup.policy
# ============================================================================
hr "정답 3"
run_app count
produce s13_orders <<'EOF'
C001:{"order_id":"O-2001","amount":1000}
C002:{"order_id":"O-2002","amount":2000}
EOF
sleep 6

# (a) 확인 명령
KT --describe --topic s13-count-app-order-counts-changelog
stop_apps

# 기대 출력:
#   Topic: s13-count-app-order-counts-changelog  ... Configs: cleanup.policy=compact,...
#
# 해설:
#   (b) delete 였다면?
#     changelog 는 상태 저장소의 "진짜 원본"입니다. 앱을 재기동하면(컨테이너 재생성,
#     새 인스턴스 투입, 로컬 디스크 유실) 이 토픽을 처음부터 끝까지 읽어 상태를 복원합니다.
#     delete 정책이면 retention.ms(기본 7일)가 지난 레코드가 삭제됩니다.
#     → 복원할 때 "7일 전에 마지막으로 갱신된 키"의 값이 사라져 있습니다.
#     → 그 키의 카운트가 0 부터 다시 시작합니다.
#     ★ 에러 없이, 조용히, 일부 키의 집계만 틀립니다. 발견하기 매우 어렵습니다.
#
#     compact 는 "키별 최신 값은 영원히 보존"이므로 이 문제가 없습니다.
#     C001 → 1, 2, 3 중 3 만 남아도 상태 복원에는 아무 지장이 없습니다.
#     Step 09 의 압축 토픽 설계 기준이 그대로 적용된 사례입니다.
#     Connect 의 connect-configs/offsets/status 도 같은 이유로 compact 입니다.
#
#   (c) repartition 은 왜 delete 여도 됩니까?
#     repartition 토픽은 "잠깐 거쳐 가는 통로"일 뿐입니다. 상태가 아닙니다.
#     Streams 가 소비 직후 스스로 purge(deleteRecords)합니다.
#     repartition.purge.interval.ms(기본 30초)마다 도는 이 로직 덕분에
#     retention.ms=-1(무한)로 설정돼 있어도 데이터가 안 쌓입니다.
#     지워져도 문제없습니다. 원본 입력 토픽에서 다시 만들 수 있으니까요.
#
#     비교:
#       repartition 을 지우면 → Streams 가 다시 만들고 데이터도 재생성됩니다 (무해)
#       changelog 를 지우면   → 상태를 영구히 잃습니다 (치명적)


# ============================================================================
# 정답 4 — 앱을 완전히 초기화하기
# ============================================================================
hr "정답 4 — 준비"
run_app count
produce s13_orders <<'EOF'
C001:{"order_id":"O-3001","amount":1000}
C001:{"order_id":"O-3002","amount":2000}
C001:{"order_id":"O-3003","amount":3000}
C001:{"order_id":"O-3004","amount":4000}
EOF
sleep 6
stop_apps

hr "정답 4 — 초기화 (세 줄)"

# ① 앱 종료 (이미 했습니다)
docker exec kafka-1 pkill -f 'Practice' >/dev/null 2>&1 || true
sleep 5

# ② Kafka 쪽 정리 — 오프셋 리셋 + 내부 토픽 삭제
K kafka-streams-application-reset.sh --bootstrap-server "$BS" \
  --application-id s13-count-app --input-topics s13_orders

# ③ ★ 로컬 상태 디렉터리 삭제 — reset 도구는 이걸 안 지웁니다
docker exec kafka-1 rm -rf /tmp/kafka-streams/s13-count-app

echo "--- 재기동. 카운트가 1 부터 시작해야 합니다. ---"
run_app count
sleep 5
stop_apps

# 해설:
#   ★ ③ 이 이 문제의 핵심입니다.
#
#   kafka-streams-application-reset.sh 는 "Kafka 쪽"만 정리합니다:
#     - 입력 토픽의 컨슈머 그룹 오프셋을 0 으로
#     - 내부 토픽(repartition/changelog) 삭제
#
#   각 인스턴스의 state.dir(기본 /tmp/kafka-streams/<app-id>)에 있는
#   RocksDB 파일은 그대로 남습니다. reset 도구는 그 머신에 접근할 수 없으니까요.
#
#   ③ 을 빼면 이렇게 됩니다:
#     - 오프셋은 0 → 입력을 처음부터 다시 읽습니다
#     - 로컬 RocksDB 에는 옛 카운트(4)가 그대로 있습니다
#     - 옛 상태 위에 새로 읽은 4건이 얹힙니다 → 카운트가 8 이 됩니다
#   ★ 에러는 없습니다. 숫자만 두 배가 됩니다.
#
#   앱이 살아 있는 상태에서 reset 을 돌리면 이렇게 거부됩니다:
#     ERROR: Java class 'kafka.tools.StreamsResetter' failed:
#       java.lang.IllegalStateException: Consumer group 's13-count-app' is still
#       active and has following members: [s13-count-app-8f2c...-StreamThread-1-consumer].
#       Make sure to stop all running application instances before running the reset tool.
#   그래서 ① 이 먼저입니다. --force 로 강제할 수도 있지만 권하지 않습니다.
#
#   코드로 하는 방법:
#     KafkaStreams.cleanUp()  ← 로컬 state.dir 을 지웁니다
#   단 start() 전에만 호출할 수 있습니다. 그리고 무조건 부르면 재기동마다
#   changelog 전체 복원이 일어나 기동이 수십 분씩 걸립니다.
#   운영 앱에는 --reset 같은 기동 플래그를 만들어 두고 그때만 부르는 패턴이 흔합니다.
#
#   그리고 reset 도구는 컨슈머 그룹 자체를 안 지웁니다. 오프셋만 0 으로 되돌립니다.
#   그룹까지 없애려면: kcg --delete --group s13-count-app


# ============================================================================
# 정답 5 — co-partitioning
# ============================================================================
hr "정답 5"
KT --create --topic s13ex_a --partitions 2 --replication-factor 3 --if-not-exists
KT --create --topic s13ex_b --partitions 4 --replication-factor 3 --if-not-exists
KT --create --topic s13_left  --partitions 3 --replication-factor 3 --if-not-exists
KT --create --topic s13_right --partitions 6 --replication-factor 3 --if-not-exists

echo "--- (a)(b) 예외 재현 ---"
docker exec kafka-1 sh -c "cd /tmp && timeout 45 java -cp '/opt/kafka/libs/*' Practice.java copartition" 2>&1 | grep -A3 'TopologyException' | head -6 || true
stop_apps

echo "--- (c) 수정판 ---"
docker exec kafka-1 sh -c "cd /tmp && timeout 45 java -cp '/opt/kafka/libs/*' Practice.java copartition --fix" 2>&1 | head -12 || true
stop_apps

# 해설:
#   (b) 예외 메시지:
#     org.apache.kafka.streams.errors.TopologyException: Invalid topology:
#       Following topics do not have the same number of partitions:
#       [s13_left(3), s13_right(6)]
#
#     ★ 이 예외는 build() 가 아니라 "파티션 할당 시점"에 납니다.
#       즉 streams.start() 이후 첫 리밸런싱에서 터집니다.
#       StreamsPartitionAssignor.assign() 안의 verifyCopartitioning() 이 던집니다.
#       그래서 로컬에서 토폴로지만 짜 보면 아무 문제가 없어 보입니다.
#
#   (c) 고치는 한 줄:
#     .repartition(Repartitioned.<String,String>as("right-fixed").withNumberOfPartitions(3))
#
#     이러면 s13-copart-app-right-fixed-repartition 이라는 3파티션 토픽이 생기고,
#     조인은 s13_left(3) × right-fixed-repartition(3) 으로 이뤄집니다.
#     대가는 Kafka 왕복 한 번이 추가되는 것입니다.
#
#     대안: 토픽 자체의 파티션 수를 맞춥니다.
#       kt --alter --topic s13_left --partitions 6
#     ⚠️ 단 파티션을 늘리면 Step 03 에서 본 대로 "같은 키가 다른 파티션으로" 갑니다.
#       기존 데이터의 파티셔닝이 깨지므로 운영에서는 신중해야 합니다.
#       repartition 쪽이 대개 안전합니다.
#
#   (d) co-partitioning 의 세 조건 중 Streams 가 검증하는 것은 "1개"뿐입니다.
#
#     1. 파티션 수가 같을 것       ← ★ 검증됨. 위 예외로 즉시 알려 줍니다.
#     2. 키의 타입/직렬화가 같을 것 ← 검증 안 됨
#     3. 파티셔너가 같을 것         ← 검증 안 됨
#
#     ★ 2·3 이 더 위험합니다. 왜?
#       1번은 앱이 아예 안 뜹니다. 배포 즉시 알게 됩니다. 즉 "안전한 실패"입니다.
#       2·3 번은 앱이 정상적으로 뜨고, RUNNING 이고, 랙도 0 이고, 예외도 없습니다.
#       그런데 조인 결과가 비어 있거나 일부만 나옵니다.
#
#       예: 왼쪽 토픽은 키가 String "1001", 오른쪽은 Long 1001L.
#           사람 눈에는 같은 키지만 직렬화 바이트가 달라 murmur2 해시가 다르고,
#           따라서 다른 파티션에 들어갑니다. 어느 태스크도 둘을 함께 볼 수 없습니다.
#
#       예: 한쪽 토픽을 커스텀 파티셔너를 쓰는 Java 프로듀서가 채웠다.
#           파티션 수도 같고 키 타입도 같은데 배치 규칙이 달라 조인이 안 됩니다.
#
#     방어책:
#       - 조인 전에 양쪽 다 repartition() 을 태워 Streams 의 파티셔너로 통일합니다.
#       - 참조 데이터가 작으면 GlobalKTable 을 씁니다. 전체를 복제하므로
#         파티셔닝 자체가 무의미해지고 co-partitioning 요구가 사라집니다.
#       - 키 Serde 를 양쪽에 명시적으로 같게 지정합니다(StreamJoined.with(...)).


# ============================================================================
# 정답 6 — 윈도우 집계가 실제보다 작은 원인
# ============================================================================
hr "정답 6"
docker exec kafka-1 sh -c "cd /tmp && timeout 90 java -cp '/opt/kafka/libs/*' Practice.java window --inject-late" 2>&1 | tail -12 || true
stop_apps

echo "--- JMX 로 dropped-records-total 을 직접 읽는 법 (앱이 떠 있어야 합니다) ---"
cat <<'JMX'
  docker exec kafka-1 /opt/kafka/bin/kafka-run-class.sh kafka.tools.JmxTool \
    --object-name 'kafka.streams:type=stream-task-metrics,thread-id=*,task-id=*' \
    --attributes dropped-records-total \
    --jmx-url service:jmx:rmi:///jndi/rmi://localhost:9999/jmxrmi \
    --one-time true
JMX

# 해설:
#   (a) 99000 은 어디로 갔습니까?
#     버려졌습니다. 윈도우는 [10:22:00 ~ 10:23:00) 이고 grace 가 30초이므로
#     10:23:30 까지만 늦은 레코드를 받습니다. 그 시점에 스트림 시간은 이미
#     10:23:45 였으므로 윈도우가 닫힌 뒤였습니다.
#     ★ 예외도 안 던지고, 기본 로그 레벨에서는 아무것도 안 찍힙니다.
#       결과 토픽의 합계가 그냥 작습니다.
#       "Kafka 집계 결과가 DB 집계와 안 맞는다"의 가장 흔한 원인입니다.
#
#     ⚠️ 스트림 시간은 "벽시계"가 아니라 "지금까지 본 레코드의 최대 타임스탬프"입니다.
#       레코드가 안 들어오면 스트림 시간도 안 흐릅니다. 이것 때문에 저트래픽
#       토픽에서는 오히려 늦은 레코드가 잘 받아들여지기도 합니다.
#
#   (b) 지표 이름:
#     kafka.streams:type=stream-task-metrics,thread-id=<t>,task-id=<task>
#       dropped-records-total   ← 누적 개수
#       dropped-records-rate    ← 초당 비율
#     ★ 이 지표에 알람을 거는 것이 유일한 방어입니다.
#       "dropped-records-total 이 0 보다 크면 즉시 알람" 이 운영 표준입니다.
#       Step 14 의 JMX 절과 이어집니다.
#
#   (c) 고치는 방법:
#     TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30))
#                                                       ~~~~~~~~~~~~~~~~~~~~~~~
#                                                       이 값을 늘립니다
#     예: Duration.ofMinutes(5)
#
#     단 대가가 있습니다:
#       - 최종 결과가 그만큼 늦게 나옵니다 (suppress 를 쓴다면 더 체감됩니다)
#       - 윈도우 상태를 그동안 들고 있어야 해서 메모리/디스크를 더 씁니다
#
#   (d) 값을 정하는 기준:
#     ★ 관측된 최대 지연의 1.5 ~ 2배.
#     지연의 구성 요소:
#       - 프로듀서 재시도 시간 (delivery.timeout.ms, 기본 120초)
#       - 네트워크/브로커 지연
#       - 배치 지연 (linger.ms)
#       - 소스 시스템에서 Kafka 까지의 지연 (Connect 폴링 주기 등)
#     실무 절차: 처음에는 넉넉히(예: 5분) 잡고, dropped-records-total 을
#     며칠간 관찰하며 0 을 유지하는 선에서 줄여 나갑니다.
#
#   ⚠️ 3.0 마이그레이션 사고:
#     TimeWindows.of(size) 는 3.0 에서 deprecated 되었습니다.
#     옛 API 의 기본 grace 는 "24시간"이었습니다.
#     그런데 IDE 가 제안하는 대로 ofSizeWithNoGrace(size) 로 바꾸면
#     ★ grace 가 0 이 됩니다. 24시간 → 0.
#     조금만 늦게 도착해도 전부 버려집니다.
#     반드시 ofSizeAndGrace(size, grace) 를 쓰고 grace 를 명시하세요.


# ============================================================================
# 정리
# ============================================================================
hr "정리"
stop_apps
for app in s13-stateless-app s13-selectkey-app s13-count-app s13-window-app \
           s13-suppress-app s13-join-app s13-copart-app s13-eos-app; do
  K kafka-streams-application-reset.sh --bootstrap-server "$BS" \
     --application-id "$app" --input-topics s13_orders >/dev/null 2>&1 || true
done
for g in $(KCG --list | grep -E '^s13' || true); do KCG --delete --group "$g" >/dev/null 2>&1 || true; done
docker exec kafka-1 rm -rf /tmp/kafka-streams
for t in $(KT --list | grep -E '^s13' || true); do KT --delete --topic "$t" >/dev/null 2>&1 || true; done
docker exec kafka-1 rm -f /tmp/Practice.java
echo "정리 완료"