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 的取舍——代码本身只是把那些机制具象出来的载体。
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。
<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> -->
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 要先手动建好。
# 输入 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 实例把拓扑跑起来并挂上关停钩子。
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();
}
}
}
group.id,还是所有内部 topic 的命名前缀(01 §1.4 的 instance/thread/task 模型就架在这个消费组上)。DEFAULT_*_SERDE没配它、又没在算子里显式给 serde,运行时会在第一条记录上抛 ClassCastException 或序列化异常。start()在调用前拓扑只是个对象,调用后才加入消费组、触发首次再平衡、开始处理。
三个最常见的环境失败
| 症状 | 根因 | 修复 |
|---|---|---|
启动即抛序列化 / 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 只是在它上面加算子。
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(); }
}
}
逐行解读(代码 → 概念,不是代码 → 语法)
Consumed.with(...) 显式声明读入时的 key/value serde,不依赖默认值,类型一目了然。
count 依赖 key 共分区,于是 Streams 在这里自动插一条 word-count-app-counts-store-repartition topic,把数据按词重洗,保证同一个词落同一 task。Grouped.with(...) 给这条重分区 topic 指定 serde。
Materialized 把结果物化进名为 counts-store 的本地 RocksDB,Streams 同时为它建一条 compacted changelog(word-count-app-counts-store-changelog)——这条日志才是计数的真相源,崩溃后回放它重建。值类型是 Long,所以这里必须 .withValueSerde(Serdes.Long()),否则默认 String serde 会序列化失败。
toStream() 把「每次计数变化」转成 KStream 的一条条变更记录,再 to(...) 写回输出 topic。Produced.with(...) 声明写出 serde——值是 Long,对应读取方要用 Long deserializer。
运行 + 预期输出
编译打包后跑起应用,另开两个终端:一个用 console producer 往 text-lines 灌文本,一个用 console consumer 读 word-counts-output。
# 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
# 灌入 "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。
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()));
这三处留白不是语法填空,是选错就上线出 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,再 joinusers维表富集。 - 维表用
KTable还是GlobalKTable?给出理由。 - 可选:整条拓扑开
exactly_once_v2。
窗口聚合产出的 key 是 Windowed<String>,不是裸 userId。直接拿它去 join 按 userId 分区的维表会共分区不匹配 → 静默零输出。必须先 map 把 key 从 Windowed<String> 取回 userId(同时把窗口信息塞进 value),让 join 两侧的 key 和分区方式对齐。
参考实现 + 关键决策(自己写完再展开)
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 落到不同分区号,匹配全部静默丢失——零输出无异常,是最难查的一类。
把 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。
§本章 self-check
先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。
- 一个 Streams 应用骨架由哪三块构成?
application.id在底层同时扮演了哪几个角色?把它改名会发生什么? - WordCount 里
flatMapValues、groupBy、count三个算子,哪个是 stateless、哪个改了 key、哪个建了 changelog?分别对应 01/02 的哪个概念? - 什么时候必须在
Materialized/Consumed/Produced上显式给 serde、什么时候可以靠默认 serde?给一句可操作的判据。 - 窗口聚合产出的 key 是什么类型?为什么直接拿它去 join 一张按业务 key 分区的维表会静默零输出?怎么修?
答案(先做完再展开)
- 三块:① 一份
Properties(必含application.id+bootstrap.servers+ 默认 key/value serde);② 一个StreamsBuilder描述拓扑、build()出Topology;③ 一个KafkaStreams实例start()跑起来并挂关停钩子。application.id同时是底层消费组的group.id、所有内部 topic(repartition/changelog)的命名前缀、以及应用的部署身份。改名 = 一个全新消费组 + 全新内部 topic 前缀 + offset 从头,等于状态被清空重算——要重置状态该用 reset 工具而非改名。 flatMapValues是 stateless(一拆多、只改值不改 key、不需状态);groupBy改了 key(把词设成 key)→ 触发 重分区;count是 stateful,产出 KTable 并背一条 compacted changelog(真相源)。- 判据:只要某个算子处理的实际 key/value 类型与 Properties 里的默认 serde 不一致,就在那里显式给 serde;一致时可省。WordCount 里
count的值是Long、默认是String,所以Materialized必须.withValueSerde(Serdes.Long()),Produced写出也要 Long serde。类型不一致还靠默认 serde 会在序列化时抛异常(04 章的 serde 失败)。 - 窗口聚合产出的 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 给拓扑加了什么节点。