00 · 起点 / Entry

Kafka Streams:把流处理架在消费组之上

Kafka Streams 没有独立集群——它是嵌进你 Java 应用的一个库,把流处理建在消费组、本地状态和 changelog 日志这三块基石上。这套教程讲透它的机制、失败模式与取舍,目标是系统理解 + 中高级面试。

基于版本:Apache Kafka 4.3(2026-05),Java DSL · 阅读时间:约 2–3 小时 · 代码验证状态:示例基于 kafka-streams 4.x,未在本机逐一运行。

01适合谁

面向懂核心 Kafka、但没用过 Streams 的 Java 后端工程师。具体说,三条前置能力:

  • 学过主 Kafka 教程或等同基础:能讲清分区、消费组、offset、副本各自是什么——读得懂 主 Kafka 教程 里"分区日志 + 消费组"那条主线,本教程不重讲这些。
  • 能读写 Java、用过 Kafka 客户端:看得懂泛型签名如 KStream<String, Long>、lambda、链式调用;在本机用 KafkaConsumer 起过一个消费组读消息。
  • 理解有状态服务的基本难题:知道"进程挂了本地内存就没了"意味着什么,能想象一个聚合计数器要怎么在重启后恢复——本教程把这个难题的 Streams 解法讲到底。

02不适合谁

三类读者在别处能拿到更对口的资源:

  • 没碰过 Kafka 的零基础读者:这里不从分区、消费组讲起。先看 主 Kafka 教程 把"Kafka 是一条分区日志"这条主线打通,再回来——Streams 的容错、缩放、EOS 全是从那条主线推导出来的。
  • 只想要 API 速查的读者:要查 aggregate 的重载或某个配置项,直接读 官方 DSL 文档 比这套教程快——这里讲的是"为什么这么设计",不是方法字典。
  • 已经专精 Flink / Spark Streaming 的读者:你已有流处理的心智模型,要的是差异点和迁移指南,而非从概念建起。直接对照 Streams 文档 找"库 vs 独立集群""changelog 复制 vs checkpoint"这两处关键差异即可。

03读完之后你能做到什么

你能把任意一个 Streams 算子的行为追溯到"它跑在消费组上、状态是一条 changelog 日志"这条主线——面试被问"有状态算子怎么容错"时,能说出 changelog 回放 + standby 副本,而不是含糊地说"框架自动处理"。

落到可验证的能力,读完这套教程之后:

  • 区分 KStream / KTable / GlobalKTable 三者,并讲清各自背后是"追加账本 / 每 key UPSERT changelog / 每实例全分区副本",以及什么场景该选哪个。
  • 推导一个有状态 task 在实例故障后如何恢复——本地 RocksDB 怎么由 compacted changelog 重建,standby 副本在哪一步把分钟级恢复压到秒级。
  • 诊断一个 join 静默无输出的故障,定位到共分区要求(相同分区数 + 相同分区方式)被破坏,并给出修复路径。
  • 判断一个场景该用 Kafka Streams、原生 consumer、还是 Flink,讲出判据(要不要状态/窗口/join/EOS、是不是以 Kafka 为源汇、要不要独立集群)。
  • 说清 exactly_once_v2 的真实边界——它只在单个 Kafka 集群内把 offset 提交、状态写、输出 produce 绑进一个事务,不覆盖数据库写、HTTP 调用这些外部副作用。

一句话本质

Kafka Streams 没有独立集群——它是嵌进你应用的库,把流处理架在消费组之上:task=分区,缩放和再平衡都是消费组那套。有状态算子的状态存在本地 RocksDB,但真相源是一条 Kafka changelog 日志,所以"流"和"表"只是同一条日志的两种视图。

这一句承接主教程的"Kafka 是日志"。容错(changelog 回放)、缩放(task = 分区上限)、EOS(事务)全由它推导。中高级面试里能不能把 Streams 框成"消费组 + 本地状态 + changelog 日志"、而不是"又一个流处理引擎",就是区分背题和真懂的分水岭。

04现状速览(截至 2026-06)

什么定论、什么在动、什么已被取代

稳定(已 GA、可放心用):版本化状态存储(KIP-889,3.5 GA,每 key 存多版本、支持 get(key, asOf) 与正确的乱序处理)、交互式查询 IQv2(KIP-796)、外键 join(KIP-213,2.4)、EOS v2(exactly_once_v2,3.0 起默认)、机架感知 standby(KIP-925,3.2)。

近 12 个月的变化:KIP-1071 Streams 再平衡协议把 task 分配从客户端移到 broker 端 group coordinator,4.2 起新集群默认开(broker + client ≥ 4.2,feature streams.version=1);它仍不支持 static membership、在线 classic→streams 迁移、运行中改 topology。最新 4.3(2026-05)加 KIP-1035(StateStore 自己管 changelog offset,恢复更快)。

已淘汰:eos-v1(exactly_once)→ exactly_once_v2;IQv1 → IQv2;streams-scala 4.3 已弃用、5.0 移除(Java DSL 不受影响)。

运行要求:Streams + Clients 需 Java 11,Broker 需 Java 17(4.0 起)。

05读之前:三个假象

流畅感不等于学会

Streams 的 DSL 读起来像在写集合操作(map / filter / groupBy / count),正因为眼熟,三种"感觉良好"会骗过你——它们都是假象:

· "我读得很顺"——顺,多半是因为 DSL 长得像 Java Stream,链式调用一眼能扫过去。但能读懂 groupByKey().count() 这行,不等于能说出这背后建了一个 RocksDB、一条 changelog topic、并把状态绑进了消费组。

· "我做题很快"——快,多半是碰上了套路题("KStream 和 KTable 区别?追加 vs 更新")。换成"两个 topic join 为什么零输出、还不报错"这种,速度立刻说明不了理解。

· "我没卡壳"——没卡壳,多半是还没碰到真正的 schema:把 Streams 当"另一个流处理引擎"时一路通畅,直到遇到"它凭什么没有独立集群""状态在本地 RocksDB 那进程挂了怎么办"才会卡——那一卡,才是开始学的地方。

06概念地图

这张图是后面六章挂载细节的骨架。中心是一个嵌入式库——它没有自己的集群,其余所有角色都围绕这个事实展开。特别留意那条朱红副线:状态虽在本地 RocksDB,真相源却是 Kafka 里的一条 changelog 日志。

Streams 应用 = 嵌入式库 KStream / KTable / GlobalKTable 流、表两种视图 流表对偶 Topology 拓扑 source→processor→sink 编译为 Task = Partition 并行单位,缩放上限 Window / Time event-time + stream-time 有状态算子用 EOS v2(事务) read-process-write 原子 端到端 State Store 本地 RocksDB 有状态算子用 Changelog topic 真相源在 Kafka 持久化到
图 0.1Kafka Streams 的概念地图——所有角色都挂在"嵌入式库"这个中心事实上。 注意三点:① 没有独立集群,库就跑在消费组上;② 状态在本地 RocksDB,但真相源是那条朱红的 changelog 日志;③ 流和表只是同一条日志的两种视图。

07学习路径建议

顶部那条 breadcrumb 是完整线性路径。但按你的目的,可以走不同子路径:

  • 吃透原理应付面试:01 概念 → 02 原理 → 06 自测。概念建词汇表,原理讲机制与取舍,自测用三层梯度题逼出"强答案必须包含什么"。跳过实操与综合不影响理解主线。
  • 动手搭一个流处理应用:01 → 02 → 03 实操 → 04 陷阱 → 05 综合。从概念到 Java DSL worked→partial→open 三阶,再过一遍生产失败模式,最后用实时订单分析把所有机制串起来。
  • 读懂别人的 Streams 代码:01 → 02 → 04 陷阱。先建概念和机制的心智模型,再直奔陷阱章——别人代码里那些 selectKey、Materialized、num.standby.replicas 的写法,多半是在绕开某个失败模式,对照陷阱章一眼看懂动机。

08目录

09学完之后

这套教程把 Streams 框成"消费组 + 本地状态 + changelog 日志"。在这个 schema 上,下一步可以往这几个方向加东西:

  • Kafka Connect 子教程——在你的 schema 上补"数据怎么进出 Kafka"这一层:source/sink 连接器把外部系统接成 topic,正好喂给 Streams 处理、再写回去。
  • ksqlDB——在 Streams 之上加一层 SQL。它把 CREATE TABLE ... AS SELECT 编译成 Streams topology,让你用 SQL 表达同样的流表对偶——理解了 Streams 再看它,等于看懂了它的底座。
  • Flink 对比——补上"独立集群 + checkpoint"这一套对照物:Flink 是有调度器的独立集群、用 checkpoint 容错,与 Streams 的"库 + changelog 复制"形成最清晰的取舍对比,是面试 discrimination 题的高频考点。
  • Schema Registry——在 serde 这一层加"契约管理":Avro/Protobuf + 兼容模式解决 schema 漂移,正好补上陷阱章里 serde 崩溃那条的系统性解法。