Chapter 03

实操:把原理写成一个能跑的 Streams 应用

02 章把六个机制拆开讲透了——状态怎么存怎么恢复、改 key 何时偷偷建 repartition topic、窗口靠 stream-time 推进、join 的共分区要求、EOS 用一个事务包住什么。理解了机制,这章动手把它写成代码:一个完整可运行的 Streams 应用,从 StreamsBuilder 到 KafkaStreams.start(),再到窗口聚合 + 维表 join。三阶递进——完整示例逐行读、半成品填决策点、开放练习自己写。

代码验证状态

本章代码基于 kafka-streams 4.x Java DSL,未在本机逐一运行。API 形态对照 4.x javadoc 写就;运行命令与预期输出按单机 KRaft broker 的行为描述。把它当可读的参照实现,落到你自己的环境时以本机编译/运行结果为准。

本章你将建立的 schema

  • 一个 Streams 应用的骨架:Properties(application.id + bootstrap.servers + 默认 serde)→ StreamsBuilder → Topology → KafkaStreams.start()
  • worked example:单词计数完整可运行——stream→flatMapValues→groupBy→count→toStream→to,每行映射回 01/02
  • partial example:窗口聚合留三个决策点(窗口类型与 grace、groupByKey vs groupBy、Materialized 显式 serde)
  • open exercise:5 分钟滚动窗口金额聚合 + KTable 维表 join 富集,自己写
  • 常见环境失败:默认 serde 没配、输入 topic 没建、application.id 改名导致状态重置

读这章前请带着 02 的结论:改 key 会触发重分区、状态真相在 changelog、stream-time 只在有数据时前进。下面每段代码都把这些结论落成一行行 DSL;逐行解读会反复指回 01 的概念和 02 的取舍——代码本身只是把那些机制具象出来的载体。

类型 算子 KStream <String,Order> stream() KGroupedStream <String,Order> groupByKey() TimeWindowed KStream windowedBy() KTable <Windowed,Long> count()/aggregate() toStream() → to(sink)
图 3.1一条窗口聚合在 DSL 里走过的类型链。注意:类型在每个算子后变了——groupByKey 把 KStream 变成 KGroupedStream、windowedBy 再变成 TimeWindowedKStream、聚合产出的是 key 被 Windowed 包过的 KTable。能在脑子里跟住这条类型链,就不会在编译期被泛型卡住。

3.1环境:依赖、broker、应用骨架

一个 Streams 应用就是一个普通 Java 进程,靠一份 Properties 接上 broker,靠 StreamsBuilder 描述拓扑,靠 KafkaStreams.start() 跑起来。

Maven 坐标

Streams 是一个库,加一个依赖就够——它会把 kafka-clients 作为传递依赖拉进来。运行 Streams + clients 需要 Java 11,broker(4.0 起)需要 Java 17。

pom.xml XML
<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-streams</artifactId>
  <version>4.3.0</version>
</dependency>
<!-- 序列化常用 JSON 时再加(本章 worked example 只用内置 Serdes,可不加):
<dependency>
  <groupId>com.fasterxml.jackson.core</groupId>
  <artifactId>jackson-databind</artifactId>
  <version>2.17.0</version>
</dependency> -->
kafka-streams这一个 artifact 就是 Streams 的全部入口;StreamsBuilder、KStream、KafkaStreams、Serdes 都在里面。streams-scala 在 4.3 已弃用、5.0 移除,Java DSL 是长期选项。

本地 broker

Streams 需要一个能连的 Kafka 集群。本地起一个单机 KRaft broker(4.x 已彻底去掉 ZooKeeper)即可——具体步骤见主教程 §3.1 的单机 KRaft 启动,这里只取结论:broker 监听在 localhost:9092,用 kafka-topics.sh 建 topic。Streams 应用不会替你建输入 topic(只会建内部的 repartition / changelog topic),所以输入 topic 要先手动建好。

create-topics.sh Bash
# 输入 topic:单词计数的文本行(3 分区——决定了并行上限)
bin/kafka-topics.sh --create --topic text-lines \
  --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1

# 输出 topic:可不预建(to() 写入时若 auto.create.topics.enable=true 会自动建),
# 但生产环境建议显式建,控制分区数与配置
bin/kafka-topics.sh --create --topic word-counts-output \
  --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1

# 验证已建:
bin/kafka-topics.sh --list --bootstrap-server localhost:9092

应用骨架:Properties → Topology → start()

每个 Streams 应用都是这同一副骨架。三块缺一不可:一份 Properties 告诉运行时连哪个 broker、用什么 application.id、默认怎么序列化;一个 StreamsBuilder 描述拓扑;一个 KafkaStreams 实例把拓扑跑起来并挂上关停钩子。

AppSkeleton.java Java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class AppSkeleton {
    public static void main(String[] args) {
        Properties props = new Properties();
        // application.id 同时就是消费组 group.id —— 改名 = 一个全新消费组 = 状态从头来
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // 默认 serde:没有在算子里显式给 Consumed/Produced/Materialized 时,用这两个兜底
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();
        // …在这里用 builder 描述拓扑(见 §3.2)…
        Topology topology = builder.build();
        System.out.println(topology.describe());   // 打印拓扑:能看清 Streams 实际建了哪些内部 topic

        KafkaStreams streams = new KafkaStreams(topology, props);
        CountDownLatch latch = new CountDownLatch(1);
        // 优雅关停:close() 会提交 offset、刷写状态、退出消费组,避免下次启动多走一次恢复
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            streams.close();
            latch.countDown();
        }));

        streams.start();   // 此刻才真正加入消费组、开始拉取与处理
        try {
            latch.await();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}
application.id这一个字符串既是 Streams 应用的身份,也是底层消费组的 group.id,还是所有内部 topic 的命名前缀(01 §1.4 的 instance/thread/task 模型就架在这个消费组上)。DEFAULT_*_SERDE没配它、又没在算子里显式给 serde,运行时会在第一条记录上抛 ClassCastException 或序列化异常。start()在调用前拓扑只是个对象,调用后才加入消费组、触发首次再平衡、开始处理。

三个最常见的环境失败

表 3.1 · 启动期最常见的三类失败:症状 → 根因
症状根因修复
启动即抛序列化 / ClassCastException 没配 DEFAULT_*_SERDE,又没在 Consumed/Produced/Materialized 里显式给 serde 设默认 serde,或在每个 I/O 算子上显式传 serde(值类型不一致时必须显式,见 04 章)
启动后卡住、报 topic 不存在或一直再平衡 输入 topic 没预先建——Streams 只建内部 topic,不建你的输入 topic 先用 kafka-topics.sh --create 建好输入 topic 再启动
改了 application.id 后状态像被清空、从头重算 application.id 即 group.id;改名等于一个全新消费组,新的内部 topic 前缀、新的 offset 起点 把 application.id 当不可变的部署身份;要重置状态用 kafka-streams-application-reset.sh 而非改名

3.2Worked example:完整可运行的单词计数

把一行行文本拆成词、按词分组、计数、写回——经典的 WordCount 是把 01 的 KStream/KTable 和 02 的重分区、状态存储一次性具象出来的最短拓扑。

为什么从这个例子起步

单词计数同时踩中四个要点:stream() 产出 KStream(每行一条事件)、flatMapValues 是 stateless 的一拆多、groupBy 改 key 触发 重分区、count 是 stateful 且产出 KTable 并背一条 changelog。读懂这一个拓扑,后面窗口和 join 只是在它上面加算子。

WordCountApp.java Java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.KeyValueStore;
import java.util.Arrays;
import java.util.Locale;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class WordCountApp {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> lines =
            builder.stream("text-lines", Consumed.with(Serdes.String(), Serdes.String()));

        KTable<String, Long> counts = lines
            // stateless:一行拆成多个词,值变了、key 还是原 key(此处不重分区)
            .flatMapValues(line -> Arrays.asList(line.toLowerCase(Locale.ROOT).split("\\W+")))
            // 改 key:把词本身设成 key —— 触发下游重分区(§2.3)
            .groupBy((key, word) -> word, Grouped.with(Serdes.String(), Serdes.String()))
            // stateful:计数,物化进名为 counts-store 的 RocksDB + 一条 changelog
            .count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store")
                       .withKeySerde(Serdes.String())
                       .withValueSerde(Serdes.Long()));

        // KTable → 变更流 → 写回输出 topic(值是 Long,需显式 Long serde)
        counts.toStream()
              .to("word-counts-output", Produced.with(Serdes.String(), Serdes.Long()));

        Topology topology = builder.build();
        System.out.println(topology.describe());

        KafkaStreams streams = new KafkaStreams(topology, props);
        CountDownLatch latch = new CountDownLatch(1);
        Runtime.getRuntime().addShutdownHook(new Thread(() -> { streams.close(); latch.countDown(); }));
        streams.start();
        try { latch.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
    }
}

逐行解读(代码 → 概念,不是代码 → 语法)

builder.stream("text-lines", …)把输入 topic 读成一个 KStream——每行文本是一条独立事件、全部保留。Consumed.with(...) 显式声明读入时的 key/value serde,不依赖默认值,类型一目了然。
.flatMapValues(...)stateless 算子:把一行拆成多个词,只改值不改 key。它不需要状态、不触发重分区——一条进、多条出。
.groupBy((key, word) -> word, …)这一步把 key 从「原 key」改成「词本身」。改 key 就置了「需重分区」标志:下游的 count 依赖 key 共分区,于是 Streams 在这里自动插一条 word-count-app-counts-store-repartition topic,把数据按词重洗,保证同一个词落同一 task。Grouped.with(...) 给这条重分区 topic 指定 serde。
.count(Materialized.as("counts-store")...)stateful 算子,产出一个 KTable(每个词→当前计数)。Materialized 把结果物化进名为 counts-store 的本地 RocksDB,Streams 同时为它建一条 compacted changelog(word-count-app-counts-store-changelog)——这条日志才是计数的真相源,崩溃后回放它重建。值类型是 Long,所以这里必须 .withValueSerde(Serdes.Long()),否则默认 String serde 会序列化失败。
.toStream().to(...)KTable 是一条变更日志(01 §1.2 的对偶):toStream() 把「每次计数变化」转成 KStream 的一条条变更记录,再 to(...) 写回输出 topic。Produced.with(...) 声明写出 serde——值是 Long,对应读取方要用 Long deserializer。

运行 + 预期输出

编译打包后跑起应用,另开两个终端:一个用 console producer 往 text-lines 灌文本,一个用 console consumer 读 word-counts-output。

run.sh Bash
# 1) 启动应用(假设已 mvn package 出可执行 jar / 或在 IDE 里 run main)
java -cp target/streams-demo-1.0.jar WordCountApp

# 2) 终端 A:输入文本行
bin/kafka-console-producer.sh --topic text-lines --bootstrap-server localhost:9092
> kafka streams kafka
> streams streams

# 3) 终端 B:读输出(值是 Long,要指定反序列化器)
bin/kafka-console-consumer.sh --topic word-counts-output \
  --bootstrap-server localhost:9092 --from-beginning \
  --property print.key=true --property key.separator=" => " \
  --value-deserializer org.apache.kafka.common.serialization.LongDeserializer
expected-output.txt Text
# 灌入 "kafka streams kafka" 后:
kafka   => 1
streams => 1
kafka   => 2
# 灌入 "streams streams" 后:
streams => 2
streams => 3

注意输出是变更流而非「每个词只发一行最终值」:每来一条改变某个词计数的记录,就发一条新的当前值。kafka 出现两次,下游就看到 kafka=>1 再 kafka=>2——这正是 01 §1.2「KTable 对外发的是 changelog」那道预测题在运行时的样子。要「每窗只发一个最终结果」得用 suppress,那是 §3.3 的决策点。

想一想

把 .groupBy((key, word) -> word, …) 换成 .selectKey((key, word) -> word) 后面什么都不接(不接 count,也不接任何聚合/join),直接 .to("out")。Streams 会建出那条 repartition topic 吗?

展开答案(先停 10 秒再点)

不会。selectKey 是 lazy 的——它只置「需重分区」标志,自己不物化任何 topic。只有当下游出现真正依赖 key 共分区的 stateful 算子(count/aggregate/join)时,Streams 才会把那条 repartition topic 建出来。后面直接 to() 写出、没有共分区需求,标志就没人来兑现。

用 topology.describe() 的输出验证最直接:有重分区时拓扑里会出现一个 Sink → repartition topic → Source 的回环,把拓扑切成两个 sub-topology;没有时就是一条直线。这就是为什么排查「凭空多出来的内部 topic」第一步永远是看 describe()。

3.3Partial example:窗口聚合,留三个决策点

把单词计数升级成「每 5 分钟每个词出现几次」,框架给全,三个关键决策留给你填——它们各对应一个会上线翻车的取舍。

需求变成带时间维度:不是「kafka 总共出现过几次」,而是「kafka 在每个 5 分钟滚动窗口里出现几次」。拓扑骨架和 WordCount 几乎一样,区别只在 groupBy 之后插一个 windowedBy(...),聚合结果的 key 从 String 变成 Windowed<String>。下面的代码把不需要思考的部分写死,把三个会影响正确性的决策留成 TODO。

WindowedWordCount.java(含 TODO) Java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import java.time.Duration;
import java.util.Arrays;
import java.util.Locale;

StreamsBuilder builder = new StreamsBuilder();

KStream<String, String> lines =
    builder.stream("text-lines", Consumed.with(Serdes.String(), Serdes.String()));

KTable<Windowed<String>, Long> windowedCounts = lines
    .flatMapValues(line -> Arrays.asList(line.toLowerCase(Locale.ROOT).split("\\W+")))

    // ── 决策点 1:groupByKey 还是 groupBy?──────────────────────────
    // flatMapValues 没改 key(key 仍是原始行 key),但聚合要按「词」分组。
    // TODO: 选 groupByKey() 还是 groupBy((k, word) -> word, ...)?
    .???(/* ... */)

    // ── 决策点 2:窗口类型 + grace 宽限期 ──────────────────────────
    // 需求是「每 5 分钟一个不重叠的窗口」。
    // TODO: 用哪种 TimeWindows?要不要 grace?grace 给多久?
    .windowedBy(/* TODO */)

    // ── 决策点 3:Materialized 要不要显式 serde?──────────────────
    // 聚合值类型是 Long。
    // TODO: 这里能只写 Materialized.as("win-counts") 吗?还是必须补 serde?
    .count(/* TODO */);

windowedCounts.toStream()
    // Windowed<String> 的 key 需要 WindowedSerdes,值是 Long
    .to("windowed-counts-output",
        Produced.with(WindowedSerdes.timeWindowedSerdeFrom(String.class, 5 * 60 * 1000L),
                      Serdes.Long()));
洞察 · 三个 TODO 各卡住一个真实取舍

这三处留白不是语法填空,是选错就上线出 bug 的决策:决策点 1 选错多付一条 repartition topic 的成本或干脆 join/聚合错乱;决策点 2 的 grace 决定迟到事件被采纳还是静默丢弃;决策点 3 漏了 serde 直接序列化崩溃。逐个想清楚再展开。

决策点 1 答案 + 决策路径(先自己定再展开)

这里要用 groupBy((k, word) -> word, Grouped.with(Serdes.String(), Serdes.String()))。

决策路径:先问「当前 key 是不是已经是目标聚合维度?」。flatMapValues 只改了值、没碰 key,此刻 key 还是输入行的原始 key(甚至是 null),而聚合要按词分组——维度不匹配,必须改 key。改 key 只能用 groupBy(它接一个 KeyValueMapper 重新指定 key),groupByKey 是「key 已经对了、直接按现有 key 分组、不重分区」的快捷路径。

反过来说,02 §2.3 的「优先 groupByKey」指的是:当 key 已经是目标维度时别画蛇添足地 selectKey 成同样的 key 再 groupBy,那会凭空多一条 repartition topic。这里 key 本就不对,groupBy 的那次重分区是必要成本,不是浪费。

决策点 2 答案 + 决策路径(先自己定再展开)

用 TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1))。

决策路径:① 「5 分钟不重叠」= tumbling 滚动窗口,用 TimeWindows(hopping 会传第二个 advance 参数造成重叠,sliding/session 是另外的类,都不符)。② grace 宽限期决定「窗口结束后,迟到的事件还能不能更新这个窗口」。4.x 推荐用 ofSizeAndGrace(...) 显式声明 grace,而不是用已弃用的 ofSizeWithNoGrace 之外的旧默认(旧版曾静默给 24h grace,是迟到事件相关问题的来源)。给多少看业务能容忍多迟的数据——给 1 分钟意味着窗口结束后 1 分钟内到达的迟到事件仍计入,过了就按 02 的规则静默丢弃。

连带约束:changelog / 窗口存储的 retention 必须 ≥ 窗口大小 + grace,否则窗口还没关、状态先被压缩清掉。

决策点 3 答案 + 决策路径(先自己定再展开)

必须显式给值 serde:.count(Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("win-counts").withValueSerde(Serdes.Long()))。

决策路径:默认 serde 在 Properties 里配的是 String(key 和 value 都是)。窗口计数的值类型是 Long,与默认值 serde 不一致——一旦不一致就不能依赖默认,否则在写状态存储/changelog 时按 String 序列化 Long 会抛异常。key serde 这里可省(窗口 key 的内层类型 String 与默认一致,Streams 会用 WindowedSerdes 包装它),但值 serde 必须补。判据:只要某个算子处的实际类型和默认 serde 不同,就在那里显式声明。

3.4开放练习:窗口金额聚合 + 维表 join 富集

把订单流按用户做 5 分钟滚动窗口金额聚合,再 join 一张用户维表,给结果加上用户名/等级——这是把本章三阶能力合起来用的一道。

需求

输入 orders(key=订单号,value=含 userId 和 amount 的订单)。要算「每个用户在每个 5 分钟滚动窗口里的下单总金额」,并用一张 users 维表(key=userId,value=用户名+等级)把结果富集成「用户名 + 窗口 + 总金额 + 等级」,写到 user-spending。这道题要用到 01 的至少四个概念——KStream / KTable、stateful 聚合、流表对偶(toStream)——并踩中 02 的一个权衡:KStream-KTable join 的共分区要求。

要求

  • 用 selectKey 把订单流改成按 userId 分组(订单原本按订单号分区,维度不对)。
  • 用 TimeWindows.ofSizeAndGrace 做 5 分钟滚动窗口 + 一个你定的 grace。
  • 聚合用 aggregate(...) 累加金额(不是 count,因为要的是金额求和),Materialized 显式给 serde。
  • 把窗口结果 toStream() 后,去掉窗口包装、把 key 还原成 userId,再 join users 维表富集。
  • 维表用 KTable 还是 GlobalKTable?给出理由。
  • 可选:整条拓扑开 exactly_once_v2。
陷阱预警

窗口聚合产出的 key 是 Windowed<String>,不是裸 userId。直接拿它去 join 按 userId 分区的维表会共分区不匹配 → 静默零输出。必须先 map 把 key 从 Windowed<String> 取回 userId(同时把窗口信息塞进 value),让 join 两侧的 key 和分区方式对齐。

参考实现 + 关键决策(自己写完再展开)
OrderSpendingApp.java(参考) Java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import java.time.Duration;

Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-spending-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 可选 EOS:把 offset 提交 + changelog 写 + 输出 produce 绑进一个事务(§2.6)
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);

StreamsBuilder builder = new StreamsBuilder();

// 维表:小而慢变 → 用 KTable(理由见下)
KTable<String, String> users =
    builder.table("users", Consumed.with(Serdes.String(), Serdes.String()));

// 订单流(假设 value 已是 "userId:amount" 形式的简化字符串;真实场景用 JSON serde)
KStream<String, String> orders =
    builder.stream("orders", Consumed.with(Serdes.String(), Serdes.String()));

KTable<Windowed<String>, Double> perUserWindow = orders
    // 改 key:订单号 → userId(聚合维度对齐;触发重分区)
    .selectKey((orderId, v) -> v.split(":")[0])
    .mapValues(v -> Double.parseDouble(v.split(":")[1]))
    .groupByKey(Grouped.with(Serdes.String(), Serdes.Double()))   // 此刻 key 已是 userId
    .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1)))
    .aggregate(
        () -> 0.0,                                  // 初始值
        (userId, amount, sum) -> sum + amount,      // 累加金额
        Materialized.<String, Double, WindowStore<Bytes, byte[]>>as("spend-store")
            .withValueSerde(Serdes.Double()));

// 关键:把窗口结果的 key 从 Windowed<String> 还原成裸 userId,才能与 users 共分区 join
KStream<String, String> enriched = perUserWindow.toStream()
    .map((winKey, sum) -> KeyValue.pair(
            winKey.key(),                           // 取回 userId
            winKey.window().start() + "|" + sum))   // 窗口起点 + 金额塞进 value
    // KStream-KTable join:流记录探当前维表值(表更新不重触发)
    .join(users, (winValue, userInfo) -> userInfo + " | " + winValue);

enriched.to("user-spending", Produced.with(Serdes.String(), Serdes.String()));

关键决策 1 — 为什么用 KTable 而非 GlobalKTable。聚合后的流已经按 userId 重分区,和按 userId 分区的 users 维表天然共分区,KStream-KTable join 成立、每实例只持有自己那几个分区的维表、内存省。GlobalKTable 的价值是「免共分区」——当 join 左侧的 key 没法和维表对齐时才用它(代价是每实例全量加载)。这里 key 已对齐,没必要让每个实例都扛全量用户表。判据:能共分区就用 KTable,对不齐又改不动分区时才上 GlobalKTable。

关键决策 2 — 为什么聚合用 groupByKey 而非再 groupBy。selectKey 已经把 key 改成 userId 了,此刻 key 就是聚合维度,直接 groupByKey。重分区由 selectKey 的标志在 groupByKey 下游的聚合处兑现一次——只此一次。groupBy((k,v)->k, …) 把已经对的 key 再设一遍,语义上等价、不会减少重分区,却让读代码的人误以为这里改了 key。判据:key 已是目标维度时用 groupByKey 表达「key 已对」,把改 key 的意图只留给真正改 key 的 selectKey。

关键决策 3 — 窗口 key 必须先还原。窗口聚合把 key 包成 Windowed<String>,它的分区落点和裸 userId 不同。不 map 回 userId 就 join,Streams 看分区数对得上、不报错,但相同 userId 落到不同分区号,匹配全部静默丢失——零输出无异常,是最难查的一类。

洞察 · 可选 EOS 小节

把 PROCESSING_GUARANTEE_CONFIG 设成 EXACTLY_ONCE_V2,这条拓扑的「输入 offset 提交 + 窗口状态 changelog 写 + 输出 produce」就被绑进一个 Kafka 事务,下游配 read_committed 才看得到已提交结果。代价是 commit.interval.ms 默认从 30000ms 降到 100ms(更频繁提交、吞吐下降),且需要至少 3 个 broker——单机本地 broker(副本因子 1)下 EOS 起不来,本地验证时先用默认的 at_least_once,部署到 ≥3 broker 集群再开 EOS。

§3.2 worked 代码全给 §3.3 partial 留 3 个决策点 §3.4 open 框架自己搭 脚手架(教程给的):多 ───────→ 少 你承担的:少 ───────→ 多
图 3.2三阶训练的脚手架递减。注意:两条趋势是镜像的——教程给的脚手架逐阶变少,你要自己承担的逐阶变多。卡在哪一阶,就回那一阶对应的 01/02 概念,而不是往下硬抠。

§本章 self-check

先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。

  1. 一个 Streams 应用骨架由哪三块构成?application.id 在底层同时扮演了哪几个角色?把它改名会发生什么?
  2. WordCount 里 flatMapValues、groupBy、count 三个算子,哪个是 stateless、哪个改了 key、哪个建了 changelog?分别对应 01/02 的哪个概念?
  3. 什么时候必须在 Materialized / Consumed / Produced 上显式给 serde、什么时候可以靠默认 serde?给一句可操作的判据。
  4. 窗口聚合产出的 key 是什么类型?为什么直接拿它去 join 一张按业务 key 分区的维表会静默零输出?怎么修?
答案(先做完再展开)
  1. 三块:① 一份 Properties(必含 application.id + bootstrap.servers + 默认 key/value serde);② 一个 StreamsBuilder 描述拓扑、build() 出 Topology;③ 一个 KafkaStreams 实例 start() 跑起来并挂关停钩子。application.id 同时是底层消费组的 group.id、所有内部 topic(repartition/changelog)的命名前缀、以及应用的部署身份。改名 = 一个全新消费组 + 全新内部 topic 前缀 + offset 从头,等于状态被清空重算——要重置状态该用 reset 工具而非改名。
  2. flatMapValues 是 stateless(一拆多、只改值不改 key、不需状态);groupBy 改了 key(把词设成 key)→ 触发 重分区;count 是 stateful,产出 KTable 并背一条 compacted changelog(真相源)。
  3. 判据:只要某个算子处理的实际 key/value 类型与 Properties 里的默认 serde 不一致,就在那里显式给 serde;一致时可省。WordCount 里 count 的值是 Long、默认是 String,所以 Materialized 必须 .withValueSerde(Serdes.Long()),Produced 写出也要 Long serde。类型不一致还靠默认 serde 会在序列化时抛异常(04 章的 serde 失败)。
  4. 窗口聚合产出的 key 是 Windowed<K>(内含原 key + 窗口起止),不是裸 K。它的分区落点与裸业务 key 不同,所以和按业务 key 分区的维表不共分区;Streams 启动只校验分区数量、不校验分区方式,分区数对得上也不报错,但相同业务 key 落到不同分区号,匹配全部静默丢失 → 零输出无异常。修复:join 前先 map 把 key 从 Windowed<K> 还原成裸 K(窗口信息塞进 value),让两侧 key 和分区方式对齐。
进阶挑战 · 刚好够不着

把 worked 的 WordCount 改成 EOS 版 + suppress 只发最终结果

在 §3.2 的 WordCount(或 §3.3 的窗口版)基础上做两处升级,不要参考实现,自己写:① 开 exactly_once_v2——想清楚它绑住了哪三件事、为什么需要 ≥3 broker、本地单机为什么起不来。② 给窗口计数加 suppress(Suppressed.untilWindowCloses(...)),让每个窗口只发一个最终结果而不是每次计数变化都发一条。回答:suppress 的缓冲状态存活在哪、它依赖什么时钟来判断「窗口关闭」、为什么一个空闲分区会让加了 suppress 的窗口永远不发结果?untilWindowCloses 与 untilTimeLimit(..., emitEarlyWhenFull) 在「保证最终性」上有什么区别?

提示(卡住再展开)

EOS:一行 props.put(PROCESSING_GUARANTEE_CONFIG, EXACTLY_ONCE_V2),但事务状态 topic 默认副本因子 3,单 broker 建不出来;它绑的是 offset 提交 + changelog 写 + 输出 produce(§2.6)。suppress:缓冲在一个内部状态存储里(也有自己的 changelog,retention 要 ≥ 窗口 + grace);它靠 stream-time 判断窗口关闭,而 stream-time 只在有记录到达时前进——空闲分区不推进 stream-time,窗口永远不关、最终结果永远不发(要保证最终性必须 untilWindowCloses + 关窗信号能到达;emitEarlyWhenFull 在内存压力下会提前发非最终结果,破坏「只发最终」的承诺)。先 topology.describe() 看 suppress 给拓扑加了什么节点。