Chapter 05
综合实战 · 实时订单分析
前四章把零件配齐了:01 章给了词汇(应用即库、KStream/KTable/GlobalKTable、流表对偶、拓扑、有状态 vs 无状态),02 章把机制讲透(状态存储+changelog、再平衡缩放、repartition、窗口与时间、join+共分区、EOS v2),03 章把 DSL 跑通,04 章演示了机制被违反时的八种症状。本章把这些零件拼成一个真实系统——承接主 Kafka 教程的订单系统,做实时订单分析。但本章的重点不是再教一遍 API,而是逼一组判别决策:同一个需求下,该用 01 章的 KTable 还是 GlobalKTable?该用 02 章的 groupByKey 还是 groupBy?该上 Streams 还是裸 consumer?做这些选择、并说清为什么,才是会用机制和真懂机制的分界线。
本章你将建立的 schema
- 一个真实流处理系统不是"写算子",而是一连串取舍:每个取舍都在前四章某条机制上落地
- 需求决定选型——吞吐/顺序/可靠性/时间语义四项需求各自钉死一组 DSL 决策
- 判别的最高层是"该不该用 Streams":状态/窗口/join/Kafka 内 EOS 是它的甜区,非 Kafka 源汇与复杂 CEP 是它的边界
- 验收标准要可观测:迟到订单仍计入、kill 实例后状态不丢、输出 topic 无重复——每条都对应一个具体配置或算子
本章的走法:先把需求拆成四个可量化维度(§5.1 项目背景),再针对每个维度做一次判别决策(§5.2 设计任务,五个判别点,每个回链不同章节),然后合上参考实现、自己写一遍(§5.3 自己实现),最后用反思问题检验迁移能力。两张图——一张拓扑架构图、一张选型决策树——是本章的骨架,建议先扫一眼再读正文。
5.1项目背景:实时订单分析
承接主教程的订单系统:上游一个 orders topic,每条记录是一笔下单(key=订单 id,value 含 userId / amount / ts)。分析团队要一个实时看板,回答两个问题——每个用户每分钟下了多少钱,以及这个用户是什么等级——结果落到下游 order-analytics topic 供看板和告警消费。把这句话拆成可量化的需求,才能映射到 DSL 决策。
| 维度 | 具体需求 | 钉死了哪个决策 |
|---|---|---|
| 时间语义 | 按事件时间做每用户 1 分钟滚动窗口的下单金额求和;迟到订单容忍 X=2 分钟内仍计入对应窗口 | 窗口类型 + grace 取值(回 §2.4) |
| 聚合维度 | 聚合键是 userId,但 orders 的 key 是订单 id——聚合前必须把记录重新按 userId 落到同一分区 | 改 key 方式 + 是否触发 repartition(回 §2.2/§2.3) |
| 数据丰富 | 结果要带用户等级(VIP / 普通),来自一张小而慢变的用户维表 topic;维表更新不要求重算历史窗口 | 维表用 KStream / KTable / GlobalKTable(回 §1.2 + §2.5) |
| 可靠性 | 下游告警按金额触发,同一窗口结果不能重复发也不能丢;只在 Kafka 内读写,无外部副作用 | at-least-once + 幂等 / exactly_once_v2(回 §2.6 + §4.C) |
吞吐:订单峰值按几万条/秒估,远超单消费者,需要横向扩展——但并行上限被 orders 的分区数钉死(§4.7)。顺序:聚合只要求"同一 userId 的订单进同一窗口",不要求全局有序——这正好是分区内有序能满足的。可靠性:金额求和对重复敏感(重复计入会虚高告警),对丢失也敏感,且全程在 Kafka 内,没有 DB/HTTP 这类 EOS 覆盖不到的外部副作用。时间语义:用事件时间而非处理时间,因为订单可能延迟到达(网络抖动、客户端补传),用处理时间会把一笔 19:59 的订单算进 20:01 的窗口。
这四项需求不是背景板——下一节每个判别决策都直接由其中一项推出。先把这张需求表记在手边。
orders 改 key 触发 repartition 后做窗口聚合,右路用户维表物化成 KTable 供查表 join,结果经一次 EOS 事务落到分析 topic。注意:两个朱红框是这条链上唯一两处"额外成本来源"——左中的 repartition topic 是 selectKey 改 key 换来的一次 produce→re-consume 往返;下方的 EOS 事务把 offset、changelog 写、输出 produce 绑成一个原子单位,代价是把提交间隔从 30s 压到 100ms。虚线框的 state store 提醒:窗口结果活在本地 RocksDB,但真相源是它背后的 changelog 日志。5.2设计任务:五个判别决策
这一节是本章的核心,也是整套教程的迁移训练所在。每个决策都给出备选方案、这个场景选哪个、为什么,并标注它回链前四章的哪一处。判别的关键不在记住"选 X",而在能复述"为什么这个需求排除了 Y"——换一个需求,答案就该换。先看汇总表,再逐条展开理由。
| 决策点 | 备选 | 这场景选哪个 + 为什么 |
|---|---|---|
| ① 用户维表怎么物化 | KStream / KTable / GlobalKTable (§1.2 + §2.5 共分区) |
KTable。维表语义是"每用户最新等级"=按 key UPSERT,天然是 KTable 不是 append-only 的 KStream。不选 GlobalKTable 是因为聚合侧已按 userId 重分区,与维表共分区成立,无需每实例全量副本;维表虽小,GlobalKTable 的每实例 bootstrap 延迟与全量内存是白付的成本。若维表与聚合键无法共分区,才退到 GlobalKTable 当逃生口。 |
| ② 聚合前怎么改 key | groupByKey() / groupBy((k,v)->userId)(§2.3 repartition + §4.B 正确性) |
必须 groupBy 改 key(本场景无法避免)。orders 的 key 是订单 id,而聚合键是 userId,两者不同,groupByKey 会按订单 id 聚合得到全错的结果。groupBy 会触发一次 repartition(图 5.0 的朱红框)——这是正确性的必要代价,不是浪费。能用 groupByKey 省掉 repartition 的前提是 key 本就对,本场景不满足。 |
| ③ 窗口类型 + grace 取值 | tumbling / hopping / sliding / session;grace=0 / 2min / 24h (§2.4 窗口 + §4.6′ 迟到丢失) |
tumbling 1min + grace 2min。需求是"每分钟"且窗口不重叠→滚动;grace 直接由"容忍迟到 2 分钟"定为 ofMinutes(2)。grace 不能设 0(迟到订单会被静默丢弃,金额求和偏低),也不该设 24h(旧默认,结果迟迟不终态、状态膨胀)。配套:changelog/store 的 retention ≥ 窗口+grace,否则迟到事件落在已被清理的窗口上。 |
| ④ 投递语义 | at-least-once + 消费侧幂等 / exactly_once_v2(§2.6 EOS + §4.8 EOS 延迟) |
exactly_once_v2。副作用全在 Kafka 内(读 orders → 写状态 changelog → 写分析 topic),正是 EOS v2 能原子覆盖的范围;金额聚合对重复敏感,让 Streams 把三者绑进一个事务比让每个下游各写一套去重逻辑更省心。代价是提交间隔默认降到 100ms、吞吐降 15–30%——可接受。若下游还要写外部 DB,EOS 就管不到那一步,得另上幂等 upsert(本场景没有,所以 EOS 足够)。 |
| ⑤ 选型:Streams / 裸 consumer / Flink | Kafka Streams / 原生 KafkaConsumer / Apache Flink (见 §5.2 选型判别 + 图 5.1) |
Kafka Streams。要有状态窗口聚合 + join + Kafka 内 EOS——三个都落在 Streams 甜区;源和汇都是 Kafka topic;无复杂 CEP、无跨多系统的事件模式匹配。裸 consumer 要手写状态存储、窗口、容错恢复、事务,等于重造 Streams;Flink 的独立集群 + checkpoint 在这个"源汇都在 Kafka、不需要独立集群"的场景里是过度投入。 |
决策 ⑤ 展开:为什么是 Streams 而不是裸 consumer 或 Flink
前四个决策都在"已经决定用 Streams"的前提下做。决策 ⑤ 是更高一层的判别——这个任务到底该不该用 Streams。这是面试高频的 discrimination 题,三条判别线索串起来就是图 5.1 的决策树:
- 要不要状态 / 窗口 / join? 不要(只是简单转发或过滤),就别上 Streams——一个原生
KafkaConsumer加几行map更轻。本场景三者都要,第一道门就指向 Streams。 - 要不要 Kafka 内 exactly-once? 要,且源汇都在 Kafka——Streams 的
exactly_once_v2是为这个场景造的,开一个配置就有。本场景金额聚合需要它。 - 是不是以非 Kafka 源汇为主、或需要复杂 CEP / ML? 是,才考虑 Flink——独立集群 + checkpoint + 丰富 connector + CEP 库是它的甜区。本场景源汇都是 Kafka topic、无 CEP,用 Flink 等于为不需要的能力付集群运维成本。
把这三条连起来:本任务在第一道门走"要状态/窗口/join→是",第二道门走"Kafka 内 EOS→是、且不是非 Kafka 源汇为主",落点就是 Streams。换个需求——比如要对接 S3、做跨流的复杂事件模式匹配——第三道门会走向 Flink。
决策 ① 选了 KTable 而非 GlobalKTable,前提是"聚合侧已按 userId 重分区,与维表共分区成立"。如果用户维表 topic 的分区数和 orders 不同(比如维表 4 分区、orders 12 分区),这个 KStream-KTable join 会发生什么?
展开答案(先停 10 秒再点)
等值 KStream-KTable join 要求两侧共分区——相同分区数 + 相同分区方式(§2.5)。分区数不一致时 Streams 在启动时校验失败、拓扑直接拒绝构建(不是静默无输出——那是分区方式不一致的症状,见 §4.4)。
这正是 GlobalKTable 当逃生口的场景:GlobalKTable 每实例消费维表所有分区,join 时按流记录提取 key 直接查全量本地副本,免共分区要求。代价是每实例全量加载(启动 bootstrap 延迟 + 内存),只在维表小而慢变时划算——用户等级表正好符合。所以决策 ① 的"选 KTable"是有条件的:条件一旦不成立(无法共分区),答案就翻成 GlobalKTable。这就是判别题的核心——记的是条件,不是结论。
5.3自己实现
判别决策做完了,现在把它写成代码。先不要看下面的参考实现——合上它,按上一节五个决策自己写一遍 topology,写完再展开对照。判别题的价值在自己走一遍决策路径,直接看答案等于把这一章当 API 文档读。
场景与验收 checklist
实现一个 Streams 应用:消费 orders,按 userId 做 1 分钟滚动窗口的金额求和,join 用户维表补等级,以 exactly_once_v2 把结果写到 order-analytics。完成标准不是"代码跑起来不报错",而是下面三条可观测的行为——每条都对应一个具体决策:
- 迟到容忍可观测:制造一笔事件时间落在某窗口内、但在窗口结束后 1.5 分钟才到达的订单(仍在 2 分钟 grace 内),它仍计入该窗口、结果被更新;再制造一笔晚到 3 分钟的(超过 grace),它不计入、且
dropped-records指标 +1。对应决策 ③。 - 状态不丢可观测:应用运行中
kill -9一个实例,另一实例(或重启后)从 standby/changelog 恢复,被中断窗口的累计金额不归零、不丢,恢复后继续在原值上累加。对应决策 ④ 与 02 章状态容错。 - 无重复可观测:下游用
read_committed消费order-analytics,在实例反复重启 / 再平衡的过程中,同一窗口的同一最终结果不出现两次。对应决策 ④(exactly_once_v2)。
验收第一条要可控地造迟到事件,关键是给 orders 的记录显式带事件时间戳并用自定义 TimestampExtractor 取它(否则默认取 record 的元数据时间戳,迟到就不可控)。stream-time 只在有新记录到达时前进(§2.4)——要触发窗口关闭与迟到判定,得在迟到事件后再灌一条事件时间更晚的记录把 stream-time 推过 windowEnd+grace。
参考实现(写完自己版本再展开)—— 关键 DSL 片段 + 每决策为何选 X
下面是 topology 的核心片段。每段注释回指 §5.2 的决策编号,说明"为什么是这行而不是另一种写法"。完整工程还需 03 章的 StreamsConfig 与 serde 装配,这里只给承载判别决策的骨架。
StreamsBuilder builder = new StreamsBuilder();
// ── 决策①:用户维表用 KTable(不是 KStream / GlobalKTable)──
// 维表语义 = 每用户最新等级 = UPSERT changelog,天然是 KTable。
// 聚合侧下面会按 userId 重分区,与本表共分区成立,无需 GlobalKTable 的每实例全量副本。
KTable<String, String> userTier = builder.table(
"user-dim",
Consumed.with(Serdes.String(), Serdes.String())); // key=userId, value=tier
// ── orders:source 是 KStream(append-only 下单账本,不是 UPSERT)──
KStream<String, Order> orders = builder.stream(
"orders",
Consumed.with(Serdes.String(), orderSerde)
// 决策③配套:显式事件时间提取,迟到才可控
.withTimestampExtractor(new OrderEventTimeExtractor()));
// ── 决策②:必须 groupBy 改 key(不是 groupByKey)──
// orders 的 key 是订单 id,聚合键是 userId,两者不同 → 必须改 key。
// groupBy 触发一次 repartition(图 5.0 朱红框),这是正确性的必要代价。
KTable<Windowed<String>, Double> perMinute = orders
.groupBy(
(orderId, o) -> o.userId(),
Grouped.with(Serdes.String(), orderSerde))
// ── 决策③:tumbling 1min + grace 2min(不是 grace=0 / 24h)──
// "每分钟"→滚动;"容忍迟到 2 分钟"→grace=ofMinutes(2)。
.windowedBy(TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(1), Duration.ofMinutes(2)))
.aggregate(
() -> 0.0,
(userId, o, sum) -> sum + o.amount(),
Materialized.with(Serdes.String(), Serdes.Double()));
// retention 默认 = 窗口 + grace;如另设须 >= 1min+2min,否则迟到落在已清窗口
// ── 决策①落地:KStream-KTable join 查表补等级(表更新不重触发历史窗口)──
// 先把窗口结果摊回 userId 这个普通 key,再与 userTier join。
perMinute.toStream()
.map((wKey, sum) -> KeyValue.pair(
wKey.key(),
new MinuteAgg(wKey.key(), wKey.window().start(), sum)))
.join(
userTier,
(agg, tier) -> agg.withTier(tier), // 查当前表值,符合"维表更新不重算历史"
Joined.with(Serdes.String(), aggSerde, Serdes.String()))
// ── 输出到下游分析 topic ──
.to("order-analytics",
Produced.with(Serdes.String(), enrichedSerde));
Topology topology = builder.build();
// topology.describe() 可看到自动建的 repartition topic — 决策②的成本在这里现形
# ── 决策④:exactly_once_v2(不是 at_least_once + 消费侧幂等)──
# 副作用全在 Kafka 内(读 orders → 写 changelog → 写 analytics),
# 正是 EOS v2 能原子覆盖的范围;金额聚合对重复敏感。
processing.guarantee=exactly_once_v2
# 代价提示:EOS 下 commit.interval.ms 默认变 100ms(非 EOS 是 30000ms)——
# 这才是"EOS 变慢"的真因。需 >= 3 broker。
# ── 验收第二条配套:standby 让 kill 实例后从秒级追尾恢复,而非分钟级回放 ──
num.standby.replicas=1
# ── 并行:先把每实例线程提到接近 task 数(=orders 分区数),再横向加实例 ──
num.stream.threads=4
application.id=order-analytics-app
为什么这套写法对应五个决策:builder.table 而非 globalTable(①,共分区成立);groupBy 而非 groupByKey(②,聚合键≠源 key);ofSizeAndGrace(1min, 2min) 而非 ofSizeWithNoGrace 或旧 24h 默认(③);processing.guarantee=exactly_once_v2(④,Kafka 内闭环);整套用 Streams DSL 而非裸 consumer 手写状态机或 Flink 作业(⑤,三个甜区命中、源汇都在 Kafka)。验收三条分别由 grace、standby+EOS、EOS 单独兜住。
合上教程,在纸上或 Excalidraw 里画出这个实时订单分析的 topology——只画这几个元素:orders 和用户维表两个 source、窗口聚合、join、sink。然后标出哪一步会产生 repartition topic。画完回到 图 5.0 对照——你把 repartition 标在了 selectKey/groupBy 改 key 那一步、还是错标在了别处?维表那一路你画的是单箭头查表(KStream-KTable),还是误画成了双向重触发?
§反思问题
先合上教程,把答案写在纸上或编辑器里。这几道题考的不是"选了什么",而是"为什么这个需求排除了别的选项"——直接点开答案等于把这一章当再读一遍。
- 五个判别决策里,哪个你做得最没把握?回看哪一章帮你定下了它?(多数人卡在决策 ① KTable vs GlobalKTable 或决策 ③ grace 取值——前者回 §1.2 + §2.5,后者回 §2.4 + §4.6′。)
- 把决策 ② 反过来:什么情况下聚合能用
groupByKey省掉 repartition?本场景为什么不满足那个条件? - 需求若改成"全局 Top-N 热销商品"(不分用户、要跨所有分区算出销量前 10 的商品),哪些决策要改?提示:聚合键从
userId变成商品 id 容易,难的是"全局"——见下方进阶挑战。
答案(先做完再展开)
- 没有标准答案,但合格的复述要点名"决策依赖哪条机制":KTable vs GlobalKTable 依赖共分区是否成立(成立用 KTable,不成立退 GlobalKTable);grace 依赖业务能容忍多久迟到(直接等于容忍窗口),并配套 retention ≥ 窗口+grace。能说出"换个条件答案就翻",就是真懂判别。
groupByKey免 repartition 的前提是聚合键 = 当前记录的 key(不改 key,Streams 知道数据已按该 key 共分区)。本场景orders的 key 是订单 id、聚合键是userId,二者不同,所以必须groupBy改 key、必然触发一次 repartition。若上游能让orders直接以userId为 key 生产,下游就能groupByKey省掉这次往返——这是把成本左移到生产端的取舍。- 要改的核心是全局聚合这一点。每用户窗口聚合天然分区并行(每个
userId落一个分区独立算),而"全局 Top-N"要求所有商品的销量汇到一处排序——这与"task=分区、并行上限=分区数"的模型冲突。常见做法是两阶段:先按商品 id 分区做局部计数(可并行),再把局部结果用一个固定 key(如常量)groupBy汇到单个分区做全局 Top-N 合并——这一步牺牲并行度换全局视图。决策 ②(改 key)变成"二次改 key 到常量 key",决策 ③(窗口)可能从滚动改成持续维护的 Top-N 表,决策 ①(join 维表)若不再需要可去掉。这就是全局聚合与并行度的取舍:全局序必然有一个单分区瓶颈。
把"每用户每分钟"改成"全局 Top-N 热销商品"
需求变更:不再按用户,而要一个实时榜单——每分钟窗口内销量前 10 的商品,单一榜单(不是每分区一份)。直接把聚合键从 userId 换成商品 id 只解决了"按商品",没解决"全局"——按商品 id 分区后,每个分区只看得到自己那部分商品,算不出跨分区的全局前 10。设计一条能产出单一全局榜单的拓扑,并说清它在哪一步牺牲了并行度。
提示(卡住再展开)
两阶段聚合。阶段一:groupBy(商品id) 做局部窗口计数,这一步保留分区并行。阶段二:把阶段一每条结果 map 成同一个常量 key(如 "GLOBAL"),再 groupBy 汇到单个分区,在那里维护一个 Top-N 数据结构(如有界优先队列存进状态存储)。瓶颈在阶段二——所有商品的局部计数挤进一个 task,并行度退化为 1。这是"要全局序就得有单点汇聚"的固有代价;能接受近似时,可改成每分区各出局部 Top-N、下游再合并,用一点精度换回并行。