Chapter 01
核心概念:Streams 是跑在消费组上的库
起点页给出一句话本质——Kafka Streams 没有独立集群,它是嵌进应用的库,把流处理架在消费组之上,状态是一条 changelog 日志的两种视图。这章把这句话拆成六个能在面试里说清楚的概念,每个都配一段场景走查。完整可运行代码留到 03 实操。
本章你将建立的 schema
- Streams 不是引擎,是库:instance → StreamThread → task → partition,并行上限被分区数钉死
- KStream 是 INSERT 账本,KTable 是按 key UPSERT 的 changelog,
null值是 tombstone - 流表对偶:表是 changelog 的物化,流是表的差异日志——这是状态能恢复的根
- 算子分两类:stateless 单条变换,stateful 要状态存储、改 key 会触发重分区(机制留给 02)
读这章前要带着主 Kafka 教程的三块底座:分区是并行单位、消费组按分区分活、offset 记录消费进度。Streams 不替换这些——它把这些当地基,往上盖了一层"有状态消费 + 流处理 DSL"。下面六节按依赖顺序展开,每节末尾点出它和下一个概念的接口。
1.1为什么用 Streams 而非裸 consumer
Streams 把"有状态消费"里那几件难做对的事——状态存哪、怎么容错、再平衡后怎么不重算——收进库里替你做。
设想一个需求:统计每个用户的实时下单数。用裸 KafkaConsumer 写,逻辑本身一行 map.merge(userId, 1, Integer::sum) 就够。难的是它周围的三件事:
状态存哪。那个计数 Map 在 JVM 堆里。进程一重启,计数清零。要持久化,得自己接一个外部 KV 库(Redis/RocksDB),自己处理读写一致。
怎么容错。进程崩了,内存里的计数没了。重启后从哪个 offset 续?续早了会重复计数,续晚了会漏。要做对,得把"计数状态"和"offset 提交"绑成一个原子动作——裸 consumer 没给你这个原语。
再平衡后怎么不重算。消费组加了一个实例,分区被重新分配。新拿到分区 7 的实例,它手里没有分区 7 历史累计的计数——要么从头重读整个分区重建(慢且重复输出),要么接受计数错误。
这三件事不是边角料,是有状态流处理的全部难点。Streams 的设计是把它们各给一个机制:状态进本地 state store(默认 RocksDB,§2.1 详解);容错靠把每次状态写也写进一条 changelog topic,崩溃后回放重建;不重算靠复用消费组的 offset 语义,并把 changelog 写、offset 提交、输出 produce 三者在 EOS 下绑进一个事务(§2.6 详解)。
Streams 像给 consumer 装的一套"有状态消费脚手架",而不是像 Flink 那样另起一个集群。失效边界:正因为它只是库,它没有独立的资源调度器、没有跨集群的 checkpoint barrier、运行中也不能改 topology——这些是独立引擎才有的能力。把 Streams 当"消费组 + 本地状态 + changelog"理解,是真懂;当"又一个流处理引擎"理解,会在容错和缩放问题上全答错。
场景走查:从裸 consumer 到一个 Streams 拓扑
同一个"实时下单计数",用 Streams 写出来是这样——状态、容错、再平衡那三件事一行都看不到,因为它们被库吞了:
StreamsBuilder b = new StreamsBuilder();
b.stream("orders", Consumed.with(Serdes.String(), orderSerde)) // KStream:每笔订单一条
.groupBy((k, order) -> order.userId()) // 按用户重新分组
.count(Materialized.as("order-counts")) // KTable:每用户当前计数
.toStream().to("order-counts-out"); // 变更流写回 topic
Map,被 count() 背后的 state store 接管,自动有了 changelog 撑腰。第 2 行groupBy 改了 key,库会自动插一个重分区步骤——这就是裸 consumer 里"分区 7 的实例没有分区 7 的计数"问题的根,Streams 用重分区保证同一 key 永远落同一 task(§2.3)。
四行代码背后,是这一整章要拆开的机制。与下一个概念的关系:第 1 行的 stream(...) 产出一个 KStream,第 3 行的 count() 产出一个 KTable——这两个类型是 Streams 全部 DSL 的两块基石,下一节把它们和 GlobalKTable 一起讲清楚。
1.2KStream / KTable / GlobalKTable
KStream 是只追加的事件流(每条都留),KTable 是按 key 取最新值的变更表,GlobalKTable 是每个实例都持有全部分区的只读副本。
同一条 Kafka topic,业务上有两种读法。"用户 alice 点击了按钮"是事实——发生过就永远成立,再来一条点击不会否定上一条,该累加。"用户 alice 的会员等级是 gold"是状态——后一条 (alice, platinum) 应当覆盖前一条,只有最新值有意义。把这两种语义塞进同一个抽象,聚合和 join 的行为就会自相矛盾。Streams 用两个类型把它们分开:事实用 KStream,状态用 KTable。
底层机制(比文档深一层)
KStream = INSERT 语义。每条记录都是独立事实,相同 key 的多条记录全部保留、依次处理。对 KStream 做 count,结果随每条记录单调增长。
KTable = UPSERT 语义,本质是一条 changelog。记录按 key 折叠:相同 key 的新记录覆盖旧值,下游只看到"每 key 当前值"。关键的一条是 value == null 是 tombstone(墓碑)——它表示"这个 key 被删除",而不是"这个 key 的值是空"。KTable 把上游 topic 当成 compacted 日志来读:日志里同一 key 留最新、tombstone 触发删除,正是 log compaction 的语义。这就是为什么文档说"KTable 是 changelog stream 的视图"——它读的就是一条 changelog。
GlobalKTable = 每实例全分区副本。普通 KTable 每个实例只持有它被分到的那些分区;GlobalKTable 让每个实例在启动时 eager 加载该 topic 的所有分区,得到一份完整副本。代价是每实例都扛全量数据(磁盘 + 内存 + 重启时的 bootstrap 延迟),所以它只配小而慢变的维表用(国家码、商品目录)。回报是它免去共分区要求——join 时流的 key 不必和表的分区方式对齐(§2.5)。
KStream 像 Git 的 commit 历史(每次提交都是一条不可变记录,全留着);KTable 像工作区当前快照(只反映最新状态,旧值被覆盖)。边界:Git 快照是把整个历史重放一遍算出来的,而 KTable 的"快照"是每个 key 独立折叠——它不是一个全局快照,是 key→最新值的字典。所以别把 KTable 想成"某个时刻的全表镜像",它是一条还在持续吸收变更的活日志。
场景走查:同一条 topic,两种读法
topic user-tier 依次到达三条记录:(alice, silver)、(bob, gold)、(alice, gold)。
- 读成 KStream:下游看到 三 条事件,依次 silver、gold、gold。它把"alice 升级"当成一次发生过的事实记下来。
- 读成 KTable:下游看到的表是
{alice: gold, bob: gold}——alice 的 silver 被第三条覆盖。再来一条(alice, null),表变成{bob: gold},alice 这个 key 被 tombstone 删掉。
KStream<String, String> tierEvents = b.stream("user-tier"); // 三条事件全留
KTable<String, String> tierNow = b.table("user-tier"); // 每用户最新等级
GlobalKTable<String, String> country = b.globalTable("country-codes"); // 每实例全量副本
stream 取 INSERT 视图。第 2 行table 取 UPSERT 视图,同一份字节、不同折叠规则。第 3 行globalTable 让本实例独自加载 country-codes 的所有分区——只因为它小且少变。
一个 KTable 先收到 (alice, 1),紧接着收到 (alice, 3)。订阅这个 KTable 变更流的下游算子,一共会看到几条记录?值分别是什么?
展开答案(先停 10 秒再点)
会看到 两 条变更:先是 (alice, 1),再是 (alice, 3)。KTable 是 UPSERT,但 UPSERT 指的是状态怎么折叠(表里 alice 这个 key 最终只有一个值 3),不是变更被吞掉。每次值变化都会向下游发一条变更记录,所以下游看到两条;但若把这张表物化后去查 get("alice"),拿到的是最终值 3。
这指向一个常被搞混的点:KTable 的"输出"是一条变更日志(changelog),不是"每次发一份全表"。下游看到的是增量变更序列,表的"当前值"是这些变更折叠后的结果。把这一点记牢,下一节的流表对偶就是顺理成章的。
与下一个概念的关系:上面这道题里"KTable 对外发的是变更日志、表是变更折叠的结果"——把这句话正着读和反着读,就是下一节的流表对偶。
1.3流表对偶
表是变更日志"折叠到每 key 最新值"的物化,流是表"逐次变更"的差异日志——同一条日志的两种视图。
这不是一个 API,是 Streams 容错的理论根基。1.1 节那个问题——"进程崩了,本地状态没了,怎么恢复"——的答案完全建立在流表对偶上。如果表和它的变更日志是可以互相转换的,那么只要那条日志在 Kafka 里安全存着,本地那份表(RocksDB 文件)随时可以丢、随时可以从日志重建。状态可恢复,根上就是这一条。
底层机制(比文档深一层)
两个方向都要成立,对偶才成立:
流 → 表(聚合 / 物化)。把一条变更流按 key 折叠,就得到表:相同 key 后值覆盖前值,tombstone 删除该 key。count()、aggregate() 做的就是这件事——它们把一条 KStream"播放"成一张 KTable。
表 → 流(差异 / changelog)。每当表里某个 key 的值变了,就向外发一条"这个 key 现在是这个值"的记录。这串记录就是表的 changelog。KTable.toStream() 做的就是这件事。
关键在于:Streams 不是"另外维护一条日志来记录表的变化"。那条 changelog 就是表的本体。本地 RocksDB 里的表只是这条日志在某一时刻折叠出的物化结果(一份缓存)。崩溃恢复时,新实例从 changelog topic 的头开始回放,逐条折叠,就重建出了崩溃前那张表——不需要快照、不需要 checkpoint。这是 Streams 和 Flink/Spark 在容错模型上的根本分野:那两者靠周期性 checkpoint 落盘整个状态,Streams 靠回放一条始终在 Kafka 里的 changelog。
流表对偶像数据库的"事务日志 ⇄ 数据表":表是当前态,日志是到达当前态的每一步。边界:数据库的日志主要服务于崩溃恢复和复制,平时查询走表;而在 Streams 里,日志是第一性的,表是派生缓存——changelog topic 是真相源,本地 RocksDB 随时可丢可重建。方向反过来了,这正是"没有独立集群也能容错"的支点。
与下一个概念的关系:流表对偶解释了"状态为什么能恢复"。但"恢复"这个动作发生在谁身上、什么时候触发?答案要先建立 Streams 的运行单位——instance、StreamThread、task、partition,下一节就拆这条链。
1.4应用即库:instance / thread / task / partition
一个 Streams 应用就是一个普通 Java 进程,内部起 N 个 StreamThread,每个线程跑若干 task,每个 task 钉死绑定一组输入分区。
"Streams 没有独立集群"这句话,落到运行时就是这条链。理解它,缩放、再平衡、并行上限这三个最常被面试问的问题就全解开了——它们不是 Streams 自创的调度逻辑,而是白嫖了消费组的那套机制。设计上为什么这么选:复用消费组,就免费拿到了再平衡、offset 管理、EOS 集成;代价是并行能力被分区数锁死。
底层机制(比文档深一层)
四级映射,从外到内:
- instance(实例)= 一个跑着你 Streams 代码的 JVM 进程。多开几个进程(多个 pod)就是横向扩容。
- StreamThread(流线程)= 每个实例内起
num.stream.threads个。每个 StreamThread 内部持有一个普通的 KafkaConsumer,加入消费组,group.id 就是你的application.id。这是"架在消费组之上"的字面实现。 - task(任务)= 调度的最小单位。task 数 = 各子拓扑输入 topic 的最大分区数。4 个分区的输入 topic,就是 4 个 task,编号 0–3。这个数在拓扑结构定下来的那一刻就钉死了,运行时不可改。
- partition(分区)= 每个 task 绑定一组输入分区(最简单情况下一个 task 对一个分区)。task 和它的分区是焊死的,不会在运行中拆开分给两个线程。
task 在"所有实例 × 所有线程"这个池子里分配。加一个实例,触发消费组再平衡,一部分 task 迁移过去。这里藏着那个最常考的结论:并行上限 = 最大输入分区数。task 总数被分区数钉死,所以能同时干活的线程数不可能超过 task 数;线程开得比 task 多,多出来的线程空转闲置。要再提并行度,唯一办法是给输入 topic 加分区——而那是个破坏性操作(改变 key 的分区落点,§4.3)。
task 像出租车的座位数:你能拉的客人数被座位钉死,多招的司机(线程)只能在车里坐着。边界:座位数是物理的、改不了,而分区数能加——只是加分区会改变"哪个 key 坐哪个座位"(key→分区的映射变了),所以扩容是破坏性的,不像招司机那么轻。这条边界是 04 章资源陷阱的伏笔。
场景走查:4 分区 topic,逐步加线程和实例
- 1 个实例、
num.stream.threads=1:1 个线程扛全部 4 个 task。能跑,但没并行。 - 1 个实例、
num.stream.threads=4:4 个线程各 1 task,吃满并行。 - 2 个实例、各 2 线程:共 4 线程各 1 task,和上一档算力相同,但抗单实例宕机。
- 再加到共 6 线程:仍只有 4 个 task,2 个线程闲置。要更快只能给输入 topic 加分区。
一个输入 topic 有 4 个分区。你部署了 6 个 StreamThread(比如 3 个实例各 2 线程)。同一时刻有几个线程在真正处理数据?多出来的线程在干什么?
展开答案(先停 10 秒再点)
4 个线程在干活,另外 2 个闲置(除非配了 standby,那它们会作为备用副本尾随 changelog 保温,§2.2,但不处理活动 task)。
原因:task 数 = 最大输入分区数 = 4,运行时钉死。task 不会被拆成两半分给两个线程,所以线程数超过 task 数时,多出来的线程拿不到 task,空转。并行上限 = 分区数——这是把 Streams 框成"消费组之上"的最直接推论,也是把 Streams 错当"无限缩放的引擎"时第一个翻车的地方。要提升并行度,得重分区输入 topic(破坏性操作),而不是加线程。
与下一个概念的关系:task 是按"子拓扑的输入分区"切出来的。"子拓扑"是什么、一个应用怎么被切成几个子拓扑——这要先看清 Streams 把你的 DSL 代码编译成了什么结构,也就是下一节的拓扑。
1.5拓扑:DSL 编译成的处理图
拓扑是 source → processor → sink 组成的有向图;你写的 DSL 被编译成它,子拓扑的边界就是内部 topic。
DSL 写起来像链式调用 .filter().map().groupBy().count(),但 Streams 运行时不直接"执行这些方法"。它先把整条链编译成一张静态的处理图(processor topology),再按图调度。理解拓扑,是从"会写 DSL"跨到"知道运行时建了几个内部 topic、并行单位怎么切"的那道门——后者才是排查性能和成本问题的入口。
底层机制(比文档深一层)
拓扑由三类节点构成:source processor(从一个 topic 读入)、stream processor(map/filter/aggregate 等变换)、sink processor(写出到一个 topic)。DSL 的每个算子被编译成一个或多个 processor 节点连起来。
关键的一层深度在子拓扑(sub-topology)边界。一张拓扑会被切成若干子拓扑,切点正是需要重洗数据的地方——也就是改了 key、下游又要按 key 聚合/join 的位置。在切点上,上游子拓扑把数据写到一个内部 topic,下游子拓扑再从这个 topic 读回来。这个内部 topic 就是重分区 topic(§2.3)。所以"子拓扑之间靠内部 topic 连接"和"task 数 = 子拓扑输入分区数"是同一件事的两面:每个子拓扑独立按它的输入分区数切 task。
调用 topology.describe() 能打印出这张图的文本结构,看到它被切成了几个子拓扑、各含哪些 processor。这是确认"运行时到底建了几个内部 topic"的第一手段,比对着 DSL 脑补可靠。
Topology topo = b.build();
System.out.println(topo.describe()); // 打印 source/processor/sink + 子拓扑边界
DSL 编译成拓扑,像 SQL 被优化器编译成执行计划:你写声明式的链,引擎产出一张固定的算子图来跑。边界:SQL 执行计划每次查询可以重新生成、自适应调整;而 Streams 的拓扑在应用启动时一次性定型,运行中不能改(这也是 KIP-1071 至今仍不支持运行中改 topology 的原因)。把它当"可热改的计划"会在运维升级时踩空。
与下一个概念的关系:图 1.3 里那个把拓扑切成两半的重分区 topic,是被一个"改了 key 的算子 + 下游有状态算子"触发出来的。哪些算子有状态、哪些会改 key——这就是最后一节 stateless vs stateful 的分界。
1.6stateless vs stateful 算子
stateless 算子单条记录就能算出结果;stateful 算子要跨多条记录维护状态,因此需要状态存储,且改 key 时会触发重分区。
这条分界决定了一个算子贵不贵。stateless 算子(map/filter)几乎零成本,来一条处理一条、不留痕迹。stateful 算子(count/aggregate/join/窗口)要维护状态,于是连带出本章前面所有机制:状态存储(1.1)、changelog 容错(1.3)、重分区与子拓扑切分(1.5)。看一眼 DSL 链里有没有 stateful 算子,就能预估它会建几个内部 topic、恢复会不会慢。
底层机制(比文档深一层)
stateless:map、mapValues、filter、flatMap、branch、merge、peek。每条记录的输出只取决于这条记录自己,运行时不为它们分配状态存储,也没有 changelog。
stateful:count、aggregate、reduce、各种 join、所有窗口算子。它们的输出取决于"这个 key 此前累积的状态",所以必须有一块状态存储记着累积值(§2.1 详解,本节只点到"需要它")。
两类之间有一条暗线:改 key 的操作会给数据流打上"需重分区"标志。selectKey、map(可能改 key)、flatMap、groupBy 都属此类。一旦改了 key,下游的 stateful 算子要求"同一 key 的所有记录落在同一 task"才能正确聚合——而改 key 后记录的分区落点乱了,于是 Streams 在它们之间自动插一个重分区 topic 把数据按新 key 重洗(§2.3 详解,含为什么 groupByKey 比 groupBy 省、以及 selectKey 为何是 lazy 的)。这就是图 1.3 里那个子拓扑边界的来历。
场景走查:一条链里哪些算子贵
b.stream("clicks")
.filter((k, v) -> v.valid()) // stateless:零状态,零内部 topic
.selectKey((k, v) -> v.userId()) // 改 key:打上"需重分区"标志
.groupByKey() // 不再改 key
.count(); // stateful:状态存储 + changelog + 触发重分区
filter 不留状态,纯白嫖。第 3 行selectKey 改了 key,但它自己 lazy,不立刻建 topic。第 5 行count 是 stateful,它既要状态存储+changelog,又因为上游改过 key 而真正物化出重分区 topic——这一条算子把前面三节的机制全勾起来了。
别把"改 key"和"建重分区 topic"画等号。selectKey 单独存在时不建任何 topic——它只是 lazy 地标记了需重分区,真正物化要等下游一个依赖 key 的 stateful 算子来兑现。所以判断一条链有几个内部 topic,看的是"改 key 标志 + 下游 stateful 算子"的组合,不是数 selectKey 的个数。确认实情用 topology.describe()(1.5),别脑补。
下面两条链,哪条会让 Streams 建一个内部重分区 topic?
A:stream("t").mapValues(v -> v.trim()).filter(...)
B:stream("t").selectKey((k,v) -> v.region()).groupByKey().count()
展开答案(先停 10 秒再点)
只有 B。A 全是 stateless 且不改 key:mapValues 按定义不碰 key,filter 也不碰——没有"需重分区"标志,也没有 stateful 算子,运行时一个内部 topic 都不建。
B 里 selectKey 改了 key(打标志),下游 count 是 stateful 且依赖 key——标志被兑现,Streams 插入一个 <app-id>-...-repartition 内部 topic 把数据按 region 重洗,于是拓扑被切成两个子拓扑。设计指向:想省这个往返,优先用不改 key 的 mapValues + groupByKey,而不是 map + groupBy(§2.3)。
承上启下:这六节搭好了 Streams 的词汇表——库的运行单位、两种数据视图、对偶、拓扑、算子分类。每一节都把"机制深一层"压到了一个停止点:状态存储到底怎么落盘、重分区的具体成本、再平衡时状态怎么搬、窗口和时间怎么算、四种 join、EOS 怎么实现——这些是 02 原理 的正题。
§本章 self-check
先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。
- 用一句话说清 KStream、KTable、GlobalKTable 三者的语义差别。哪个值代表"删除",它在哪个类型里有特殊含义?
- "Streams 的状态能在崩溃后恢复"——把这件事的根源追到流表对偶上,讲清为什么本地 RocksDB 文件随时可以丢。
- 一个输入 topic 6 个分区,你起了 2 个实例、每个实例
num.stream.threads=5。有几个 task?几个线程在干活?要把并行度提到 10,唯一的办法是什么、它的代价是什么? - (设计题)你要给一个"按城市统计实时订单额"的需求设计 DSL 链,源 topic 的 key 是订单号、不是城市。画出这条链会被切成几个子拓扑、为什么、各子拓扑的 task 数由谁决定。哪一步是把它从一个子拓扑变成两个的"扳机"?
答案(先做完再展开)
- KStream = INSERT 追加账本,相同 key 多条全留;KTable = 按 key UPSERT 的 changelog,每 key 只留最新值;GlobalKTable = 每个实例持有全部分区的只读副本(小维表用,免共分区)。
null值是 tombstone(删除),它的特殊含义在 KTable / GlobalKTable 里成立——表示删除该 key,而不是"值为空";在 KStream 里 null 只是一条普通记录的空值。 - 流表对偶说:表是变更日志折叠到"每 key 最新值"的物化,而那条变更日志(changelog topic)就是表的本体、存在 Kafka 里。所以本地 RocksDB 只是这条日志在某时刻折叠出的缓存——它丢了,新实例从 changelog topic 头部回放、逐条折叠,就重建出同一张表。真相源在 Kafka 的日志,不在本地文件,这是状态可恢复的根,也是 Streams 不用 checkpoint 就能容错的原因。
- task 数 = 最大输入分区数 = 6。共 10 个线程,但只有 6 个能拿到 task 干活,另外 4 个闲置(除非配 standby 做备用)。要把并行度提到 10,唯一办法是把输入 topic 重分区到 ≥10 个分区——代价是这是破坏性操作:它改变 key→分区的映射,已有状态/顺序假设会被打乱,通常要重建状态甚至停机迁移。
- 会被切成 两个子拓扑。链大致是
stream(订单).selectKey(按城市).groupByKey().aggregate(求和)。扳机是selectKey把 key 从订单号改成城市——下游aggregate是 stateful 且依赖 key,要求同城市记录落同一 task,于是 Streams 在两者之间插一个重分区 topic 把数据按城市重洗,拓扑就此一分为二。子拓扑 1 的 task 数由源 topic(订单 topic)的分区数决定;子拓扑 2 的 task 数由那个重分区 topic 的分区数决定(默认与源一致,可用Repartitioned调)。
同一份 topic,既要当流又要当表
有一个 account-events topic,每条是一次账户余额变动事件。需求 A:实时对账,要看到每一笔变动(审计用)。需求 B:随时查某账户当前余额。这两个需求在 Streams 里分别该把这条 topic 读成什么类型?把同一条 topic 同时建成 KStream 和 KTable,运行时会各自付出什么代价(提示:想想各自背后有没有状态存储、有没有 changelog)?如果 B 还要支持"查三天前某账户的余额",本章的 KTable 够用吗?
提示(卡住再展开)
A 用 KStream(INSERT 视图,每笔都留,无状态存储、无 changelog,几乎零额外成本);B 用 KTable(UPSERT 视图,要物化成 state store + changelog,付出存储和写放大的代价)。两者可以从同一条 topic 各读各的。至于"查三天前的余额"——普通 KTable 只留每 key 最新值,答不了历史时点查询;这正是 版本化状态存储(versioned state store)要解决的问题,它给每 key 存多版本 (value, timestamp)。本章不展开,留意 02/04 章会点到它属于 Kafka 3.5+ 的稳定能力。