Chapter 02

原理:状态、再平衡、窗口、Join、EOS 怎么工作

上一章把 Streams 框定为一个跑在消费组之上的库——task=分区、缩放走再平衡、状态有本地副本但真相源是一条 changelog 日志。这章把六个机制拆开看:状态到底怎么存怎么恢复、再平衡如何搬动有状态 task、repartition topic 何时被偷偷建出来、窗口靠什么时钟推进、四种 join 各自的代价、EOS 用一个事务包住了什么。每个机制都配一张备选方案对比表——看清"为什么是这个设计"比记住 API 更值钱。

本章你将建立的 schema

  • 本地 RocksDB 是读缓存,compacted changelog 才是真相源;恢复=回放,standby=保温
  • 有状态 task 在再平衡里被搬走时要重建状态——这是缩放卡顿的根因,不是 bug
  • 改 key 的算子给下游埋了一次 produce→re-consume 往返,repartition topic 是隐藏成本
  • stream-time 由记录的最大时间戳驱动,只在有数据到达时前进——空闲分区会冻结窗口
  • 四种 join 的共分区要求与逃生口(GlobalKTable / 外键 join)
  • EOS v2 用一个 Kafka 事务绑定 offset 提交+状态写+输出,边界止于 Kafka 集群

2.1状态存储:RocksDB 是缓存,changelog 是真相

有状态算子把状态写进本地 RocksDB,同时每次 put 追加进一条 compacted changelog topic——故障转移靠回放这条日志重建。

一个有状态算子(aggregate / count / reduce / 窗口 / join)需要一块跨记录存活的内存。Streams 默认给它一个嵌入式 RocksDB 实例:一个 LSM-tree 的本地键值库,落盘到 task 的状态目录,所以进程重启后目录还在、状态不丢。RocksDB 前面再挂一个小的写回缓存(record cache),把同一 key 的连续更新合并后再下沉,减少写放大。但本地盘会随实例一起消失——容器被调度走、磁盘坏掉、task 被再平衡搬到别的实例,本地 RocksDB 就没了。所以每次 put 还会同步追加一条记录到一个内部 changelog topic(命名 <application.id>-<store-name>-changelog,配置为 log-compacted——每个 key 只保留最新值 + tombstone)。这条日志才是状态的真相源:RocksDB 丢了可以从零回放 changelog 重建,changelog 丢了 RocksDB 就成了孤儿。这正是流表对偶落到工程上的样子——table 是 changelog 的物化视图,所以 table 永远可以从日志重新长出来。

实例 A(owner) 实例 B(standby) StreamThread 处理记录 本地 RocksDB 读缓存 / 落盘 put changelog topic compacted 真相源 每 put 写 尾随 standby 副本 保温的 RocksDB 新 owner 重建 回放 changelog 故障转移
图 2.1本地 RocksDB 只是缓存,每次 put 都写进 compacted changelog——后者才是真相源。注意:实例 A 挂掉后,新 owner 不是从 A 的盘恢复(盘没了),而是回放 changelog;standby 提前尾随同一条日志,把"分钟级回放"压成"秒级追尾"。

故障转移的代价由状态大小决定。一个多 GB 的状态从零回放 changelog 会阻塞启动数分钟——这是 Streams 运维里最常被点名的痛点(04 章专门讲它)。缓解手段是 standby 副本(num.standby.replicas):另一个实例后台尾随同一条 changelog,把状态"保温"在本地 RocksDB 里;owner 挂掉时再平衡把 task 优先分给持有 standby 的实例,它只需追上最后一小段日志,恢复从分钟降到秒。Kafka 4.3 的 KIP-1035 进一步让 StateStore 自己管理 changelog 的已恢复 offset,重启后能跳过已重放的部分,恢复更快。

word-count-store.java Java
// 一个有状态 count:状态进本地 RocksDB,并由内部 changelog 撑腰
KStream<String, String> lines = builder.stream("text-lines");
KTable<String, Long> counts = lines
    .flatMapValues(v -> Arrays.asList(v.toLowerCase().split("\\W+")))
    .groupBy((key, word) -> word)            // 改 key → 触发重分区(见 §2.3)
    .count(Materialized.<String, Long>as(    // 物化为一个具名状态存储
        Stores.persistentKeyValueStore("word-counts")));
// 这块状态对应的内部 topic:
//   <application.id>-word-counts-changelog   (compacted,真相源)
count(Materialized…)把聚合结果物化进名为 word-counts 的 RocksDB 存储,Streams 同时为它建一条 compacted changelog topic;这条 topic 的 retention 必须足够长,否则压缩后丢的就是真相。
表 2.1 · 状态放哪:远程状态库 vs 本地 RocksDB + changelog
方案优势为什么没选 / 选中
远程状态库(Redis / Cassandra / 外部 KV) 状态与计算解耦,实例无盘、再平衡不用搬状态 每次 get/put 走一跳网络(毫秒级 vs 本地亚毫秒),热路径吞吐塌方;还要单独运维一套 HA 集群、自己解决一致性
纯内存状态(无落盘、无 changelog) 最快,零磁盘、零外部 topic 进程一重启状态全丢,无法容错;状态超过堆就 OOM——只适合可重算的小状态
本地 RocksDB + compacted changelog 本地读亚毫秒、无逐条网络跳;changelog 复用 Kafka 已有的复制/HA,无需第二个数据库 选中:把"复制层"外包给 Kafka 自己——状态在哪算就在哪存,容错由日志兜底

带来的代价

① 恢复阻塞:多 GB 状态从零回放 changelog 会阻塞启动数分钟,无 standby 时每次再平衡都会触发。② 写量翻倍:每次 put 既写本地 RocksDB 又写一条 changelog 记录,broker 端多承一份写入和存储。③ 堆外内存:RocksDB 的 block cache 与 memtable 分配在 JVM 堆之外,-Xmx 管不住——容器内存按"堆 + RocksDB 缓存 + 余量"算,否则被 OOMKilled(04 章资源陷阱)。④ 版本化状态存储(KIP-889, 3.5 GA)给每个 key 存多版本 (value, timestamp) 以支持 get(key, asOf) 和正确的乱序处理,代价是更高的存储与写放大。

想一想

一个实例持有 8GB 状态,没有配 standby 副本。它崩溃后被 Kubernetes 用一个全新 pod(空磁盘)拉起。新 pod 要多久才能开始处理新记录?如果改用持久卷(StatefulSet)让磁盘跨重启存活,会快多少?

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

空盘启动必须从头回放整条 changelog 重建 8GB RocksDB——受 broker 读带宽和 RocksDB 写入速度限制,分钟级甚至更久,期间该 task 不处理任何新记录。这就是"再平衡卡顿数分钟"的来源。

持久卷让本地 RocksDB 跨重启存活:新 pod 挂回旧盘,状态目录还在,只需从 changelog 追上崩溃后那一小段增量,恢复降到秒级。这也是为什么 standby 副本(在别的实例保温)和持久卷(在本地保活)都能解同一个问题——它们都让"恢复"从"全量回放"退化成"追尾增量"。

2.2再平衡与缩放:消费组那套,外加搬状态的代价

缩放就是消费组再平衡——task 在实例和线程间重新分配,但有状态 task 被搬走时要在新落点重建状态。

Streams 不自造调度器。每个 JVM 实例起 num.stream.threads 个 StreamThread,每个线程内嵌一个普通 consumer,全部以 application.id 作为 group.id 加入同一个消费组。task = 子拓扑输入 topic 的每个分区一个,所以 task 总数被最大输入分区数钉死、运行时不可改。加一个实例、某实例掉线、某线程死掉,都触发一次再平衡,把所有 task 在当前所有实例的所有线程间重新摊分。无状态 task 搬家几乎零成本——换个线程接着消费就行。有状态 task 不同:它在新落点没有本地 RocksDB,必须先回放 changelog 重建状态(或从 standby 追尾)才能处理,这段时间该分区的处理停滞。所以"缩放"在有状态拓扑里从来不是瞬时的——这是机制的必然代价,不是故障。

分配逻辑在哪算,是 4.x 的关键变化。经典协议(classic)由其中一个 client 当 group leader、在客户端计算 task 分配;KIP-1071 Streams 再平衡协议把分配移到 broker 端的 group coordinator,4.2 起新集群默认开(需 broker + client 都 ≥ 4.2,feature streams.version=1)。broker 端分配能拿到全局视图做更稳的 task 放置与 standby 调度,但仍不支持 static membership、运行中在线把 classic 协议迁到 streams 协议、以及运行中修改 topology。

表 2.2 · 再平衡协议:classic(客户端分配)vs KIP-1071(broker 端分配)
方案优势为什么没选 / 选中
自造专用调度器(独立 master 分配 task) 可针对流处理做最优放置 得自己实现成员发现、故障检测、再平衡、与 EOS 事务的协同——把消费组已经解决的问题重造一遍
classic 协议(client 端 leader 分配) 复用消费组,零额外 broker 改动;老集群唯一选项 分配在客户端算,leader 只有局部视图;大消费组的 join/sync 来回开销与"停顿式"再平衡更明显
KIP-1071(broker 端 group coordinator 分配) broker 全局视图、更稳的 task/standby 放置、再平衡更平滑;4.2 新集群默认 选中(新集群):把分配收归 broker;代价是仍不支持 static membership 与在线协议迁移

带来的代价

① 有状态搬迁成本:每次再平衡,被重分配的有状态 task 都要重建状态——无 standby 时这就是分钟级停顿。② 并行上限 = 最大输入分区数:num.stream.threads 把线程提到 task 数能用满单实例,再加实例才有意义;线程/实例超过 task 数则纯闲置,要更多并行只能重分区输入 topic(破坏性操作)。③ 无 static membership(KIP-1071 下):滚动重启 / pod 漂移更容易触发全量再平衡,靠 standby 副本和持久卷缓解,而非靠 group.instance.id 跳过再平衡。

想一想

输入 topic 有 6 个分区。部署 4 个实例、每个 num.stream.threads=3,共 12 个线程。能跑多少个并行 task?多出来的线程在干什么?要再提并行度,唯一的办法是什么?

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

task 数 = 最大输入分区数 = 6,被钉死。12 个线程里只有 6 个各领一个 task,另外 6 个线程空闲(持有 consumer 但分不到 task)。加机器、加线程都越不过 6 这个天花板。

唯一提并行度的办法是把输入 topic 重分区到更多分区(比如 12),让 task 数随之上去——但这是破坏性操作(改变 key→分区映射、打乱既有顺序与状态归属),通常意味着重建拓扑。这就是为什么分区数要在设计阶段按峰值并行需求预留。

2.3重分区:改 key 就给下游埋了一次往返

改 key 的算子置一个"需重分区"标志,下游第一个 stateful 算子自动建一条 repartition topic 把数据按新 key 重洗。

聚合和等值 join 要求"相同 key 的记录落在同一个 task/分区"(共分区)。但改 key 的操作——selectKey / map / flatMap / groupBy——会让记录的新 key 与它当前所在分区不再一致。Streams 的处理方式:这些算子只是置一个"需重分区"标志,自己不立即重洗;当下游出现第一个依赖 key 的 stateful 算子(aggregate / count / join)时,Streams 在它前面自动插入一条内部 <application.id>-<name>-repartition topic:上游把改了 key 的记录 produce 进这条 topic(按新 key 分区),下游再把它重新消费进来——此时相同新 key 必然同分区,聚合/join 才正确。selectKey 是 lazy 的:单独用它(后面没有依赖 key 的算子)不会建任何 topic,只有真正需要共分区的下游才把重洗物化出来。这是"默认正确"的设计——用户不必手动保证共分区,代价是这条往返被藏在拓扑里。

上游 下游 source 流 key=A selectKey 改 key→B(lazy) produce repartition topic 按新 key 分区 re-consume aggregate 同 key 已同分区 produce→re-consume 一次额外往返
图 2.2改 key 的算子(selectKey)不立即重洗,而是给下游 aggregate 埋了一条 repartition topic——数据先 produce 出去、再 re-consume 回来。注意:这一次 broker 往返是延迟与磁盘成本的来源;selectKey 单独用不建 topic,是下游的 aggregate 把它物化出来的。
表 2.3 · 共分区怎么保证:强制用户共分区 vs 自动 repartition
方案优势为什么没选 / 选中
强制用户事先共分区(改 key 后 join 直接报错,要用户手动重发到对齐的 topic) 零隐藏 topic,成本完全显式,工程师清楚每一跳 把共分区的正确性责任全压给用户,极易出错;漏对齐就静默错或拓扑失败,DSL 的"声明式"承诺破产
每个改 key 的算子都立即重洗(不 lazy) 语义直白,所见即所建 selectKey 后若没有依赖 key 的下游,重洗纯属浪费;会凭空多出大量无用 repartition topic
改 key 置标志 + 下游 stateful 算子自动按需 repartition(lazy) 默认正确,用户不必手动共分区;只在真正需要时才物化重洗 选中:正确性默认兜底;代价是 produce→re-consume 往返被藏进拓扑,成本不显眼

带来的代价

① 一次完整往返:每个被物化的重分区点,每条记录都要 produce 进 repartition topic 再 re-consume 回来——多一份网络、broker 磁盘、序列化/反序列化和端到端延迟。② 隐藏 topic 成本:这些 topic 不在你的代码里,但占 broker 存储、计入分区配额,排查时容易遗漏。③ 能省则省:不改 key 的操作(mapValues / filter)和 groupByKey(不改 key)不触发重分区;优先它们而非 map / groupBy。用 topology.describe() 打印拓扑、数清实际建了几条 repartition topic(04 章)。

想一想

下面两段代码,哪段会建出 repartition topic?stream.selectKey((k,v)->v.userId()).to("out") 还是 stream.selectKey((k,v)->v.userId()).groupByKey().count()?

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

第一段不会。selectKey 是 lazy 的——它只置"需重分区"标志,后面接的是 to()(sink,不依赖 key 共分区),没有任何 stateful 算子来物化重洗,所以不建 topic。

第二段会。groupByKey().count() 是依赖 key 的 stateful 算子,它需要"相同 userId 落同分区",于是 Streams 在它前面插入 repartition topic,把 selectKey 改过 key 的记录重洗一遍。要点:建不建 topic 取决于下游是否真的需要共分区,不取决于你改没改 key。

2.4窗口与时间:stream-time 由数据推进,不是墙上时钟

窗口按事件时间切分,由 stream-time(每 task 见过的最大时间戳)驱动关闭——它只在有记录到达时前进。

窗口把无界流切成有界的桶来聚合,四种形状:tumbling 滚动(固定大小、不重叠,一条记录进一个窗)、hopping 跳跃(固定大小 + 更小步进 → 窗口重叠 → 一条记录进多个窗)、sliding 滑动(按记录两两之间的时间差成对落窗)、session 会话(按活动间隙合并,间隙超阈值就开新会话)。关键不在形状,而在驱动它们关闭的时钟。Streams 默认用 event-time(记录自带的时间戳,由 TimestampExtractor 提取),不是 processing-time(记录被处理的墙上时间),也不是 ingestion-time(写入 broker 的时间)。窗口由 stream-time 推进:stream-time = 这个 task 到目前为止见过的最大记录时间戳,单调不减。它只在有记录到达时前进——没有新数据,stream-time 就冻结,依赖它关闭的窗口永远不发结果。这是 event-time 换来可重放确定性(同样的输入重跑得到同样的窗口结果,与运行时刻无关)所付的代价。

迟到事件由 宽限期 grace 控制:窗口在 windowEnd 逻辑上结束后,还允许 stream-time 推进到 windowEnd + grace 之前到达的迟到记录更新该窗结果;一旦 stream-time 越过 windowEnd + grace,更晚到的记录被静默丢弃(只计入 dropped-records 指标,不报错、不进死信)。状态存储的 retention 必须 ≥ 窗口大小 + grace,否则窗口还没关、状态先被清掉。若要"每个窗口只发一个最终结果"而非每次更新都发一条中间结果,用 suppress(untilWindowCloses)(KIP-328)把中间更新压住、只在窗口关闭时发一次。

event-time → 窗口 [t0, t1) 聚合桶 grace 宽限期 t0 t1(windowEnd) t1+grace e1 纳入 迟到·grace 内 纳入 迟到·grace 外 静默丢弃 stream-time = 见过的最大时间戳,只随记录到达前进 无新数据 → stream-time 冻结 → 窗口不关
图 2.3窗口由 stream-time(见过的最大事件时间戳)推进;grace 内的迟到事件被纳入,越过 t1+grace 的被静默丢弃。注意:是"是否有更晚的记录把 stream-time 推过 t1+grace"决定迟到记录的去留,而不是墙上时钟——所以一个空闲分区会冻结 stream-time,让窗口永远不关、结果永远不发。
表 2.4 · 窗口时钟:processing-time vs event-time
方案优势为什么没选 / 选中
processing-time(按记录被处理的墙上时间分窗) 实现简单,时间永远向前、永不冻结,低延迟出结果 结果不可重放——同一份数据重跑因机器/时刻不同得到不同窗口归属;网络抖动、回填历史数据会把记录算进错误的窗
ingestion-time(按写入 broker 的时间分窗) 比 processing-time 稳定,时间在 broker 端钉死 仍不反映事件真实发生时刻;上游缓冲/批量写入会扭曲时间;回填依旧错
event-time + stream-time(按事件自带时间戳分窗) 可重放、确定性——同输入同结果,与运行时刻无关;正确处理乱序与回填 选中:正确性优先;代价是空闲/低流量分区会冻结 stream-time,窗口不关、结果不发

带来的代价

① 空闲分区冻结:stream-time 只随记录前进,低流量或空闲分区让窗口迟迟不关、结果延迟甚至永不发出——event-time 的根本代价(04 章)。② suppress 的内存压力:suppress 把未关窗口的中间状态缓在内存,缓冲满时行为取决于配置——emitEarlyWhenFull 会在内存压力下提前发出非最终结果(破坏"只发最终"的承诺),要严格"最终"必须用 shutDownWhenFull(宁可崩也不发中间结果)。③ retention 约束:窗口状态与 changelog 的 retention 必须 ≥ 窗口 + grace,配小了窗口没关状态先没。

想一想

一个 5 分钟 tumbling 窗口 + 1 分钟 grace。某分区在 10:00 收到最后一条记录后就没有任何新数据进来了。覆盖 09:55–10:00 的那个窗口,结果会在什么时候发出?

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

它不会自动发出(至少在该分区有新记录之前)。窗口要在 stream-time 推进到 windowEnd + grace(即 10:01)之后才关闭,而 stream-time = 该 task 见过的最大记录时间戳。10:00 之后没有任何记录到达,stream-time 就停在 10:00,永远到不了 10:01,窗口永远不关、结果永远不发。

这正是 event-time 的代价:墙上时间走到 10:05 也没用,Streams 不看墙钟。要让低流量场景的窗口能关,得靠业务侧持续有数据、或用 processing-time 兜底(牺牲可重放)、或周期性灌入 heartbeat 记录把 stream-time 顶上去。

2.5Join:四种语义,各自的共分区要求与逃生口

四种 join 对应四种数据形态,等值 join 要求两侧共分区,GlobalKTable 与外键 join 是免共分区的逃生口。

join 把两条流/表按 key 关联,语义由两侧是流还是表决定。KStream-KStream=窗口化 join:两侧记录都入状态存储缓冲,在窗口 + grace 内每出现一对匹配就发一条——因为两条都是无界流,必须给一个时间窗口才能界定"什么算同时发生"。KStream-KTable=查表 join:流记录到达时探查表的当前值,表自己的更新不回头重新触发已处理的流记录——典型用于流数据打维表标签。KTable-KTable=维护式 join:任一侧更新都重新计算并发出新结果,维护一个始终最新的连接视图。外键 join(KIP-213, 2.4):左表按 value 里提取的外键关联右表主键,Streams 在背后建两个隐藏 topic(subscription + response)加一个复合 RocksDB 存储来重新按外键路由——这是唯一不要求左侧用主键 join 的表-表 join。

共分区要求是等值 join 的硬约束:两侧必须相同分区数 + 相同分区方式(同一个 producer 分区器、同样的 key),相同 key 才会落到同一个 task。分区数不一致,Streams 在启动时直接校验失败;分区方式不一致(数量碰巧相同但分区器不同),Streams 不报错却静默丢匹配——这是最难排查的 join 故障。Streams 会对改过 key 的一侧自动重分区来对齐,但两个独立来源 topic 的分区方式得你自己保证。GlobalKTable 是逃生口:它在每个实例全量加载所有分区,所以任何 key 本地都能查到,免共分区,但代价是每实例一份全量副本(磁盘 + 内存 + 启动 bootstrap 延迟),只配小而慢变的维表;外键 join 同样免共分区(它自己用隐藏 topic 重路由)。

表 2.5 · 维表关联:共分区 join vs GlobalKTable join
方案优势为什么没选 / 选中
共分区的 KStream-KTable join(两侧对齐分区数 + 分区方式) 每实例只持有本分区那部分维表,内存/磁盘省;维表可大 必须保证两侧严格共分区——数量不符启动失败、方式不符静默丢匹配;维表来自别处时对齐成本高且易错
把维表 join 改成外部查询(每条流记录查一次远程 DB/缓存) 维表无需进 Streams,更新即时可见 每条记录一跳网络,吞吐塌方且引入外部依赖;丢了 Streams"本地状态"的全部优势
GlobalKTable join(维表每实例全量副本) 免共分区、任意 key 本地可查、用流记录的任意字段关联;最省心 选中(小维表):拿全量副本换掉共分区约束;代价是每实例全量内存/磁盘 + 启动 bootstrap,仅限小而慢变的维表

带来的代价

① 静默无输出:未共分区时(分区方式不一致)零输出且无任何报错——join "不工作"却查不出原因,必须核对两侧分区数 + 分区器(04 章)。② 外键 join 的内部 topic:FK join 额外建 subscription + response 两条内部 topic 加复合存储,运维面与存储成本随之上升。③ KStream-KStream 缓冲整窗:两侧记录在窗口 + grace 内都要入状态存储,宽窗 + 高流量 = 大状态、大 changelog。④ GlobalKTable 全量代价:每实例一份完整副本,维表一大就磁盘/内存爆炸、启动 bootstrap 拖慢。

想一想

两条 topic 都是 6 个分区,一个 KStream-KTable join 上线后零输出,日志里没有任何报错。分区数都对得上,问题嫌疑最大的出在哪?换成 GlobalKTable 为什么能绕过?

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

嫌疑最大的是分区方式不一致:两条 topic 由不同上游写入、用了不同的 producer 分区器(比如一边按字符串 hash、一边按自定义逻辑),导致相同 key 落到不同分区号。Streams 启动时只校验分区数量(这里都是 6,过关),不校验分区方式,于是相同 key 永远不相遇,匹配全部丢失——零输出且零报错。

GlobalKTable 绕过是因为它在每个实例全量加载所有分区:不管流记录的 key 经哪种分区方式落在哪,本地都有完整维表可查,共分区约束直接不适用。代价就是那份全量副本——所以它只对小维表成立。

2.6EOS v2:一个事务包住 offset、状态、输出

exactly_once_v2 把 offset 提交、状态 changelog 写、输出 produce 绑进一个 Kafka 事务,下游 read_committed 只读已提交。

read-process-write 循环里有三件事必须同生共死,否则就重复或丢数:消费位移的提交、状态的 changelog 写、结果的输出 produce。任意两件之间崩溃都会破坏 exactly-once——比如输出发了但 offset 没提交,重启后重新处理同一批,输出就重复。processing.guarantee=exactly_once_v2 用一个 Kafka 事务把这三者原子化:要么三件一起提交、要么一起回滚。下游消费者设 isolation.level=read_committed 就只读已提交的记录、自动跳过被回滚(aborted)的部分。v2 的关键改进(KIP-447)是一个 StreamThread 一个 producer 覆盖该线程的所有 task;v1(exactly_once,已淘汰)是一个 task 一个 producer,高分区数下 producer 数量爆炸、内存与 broker 端事务协调开销失控。这就是为什么 v2 是 3.0 起的默认、v1 被弃用。

eos-config.properties Properties
# 生产者侧(Streams 应用)
processing.guarantee=exactly_once_v2
# EOS 下 Streams 把提交间隔默认改成 100ms(非 EOS 是 30000ms)——
# 这才是"开了 EOS 就变慢"的真因,不是事务本身的固定开销
commit.interval.ms=100
# EOS 需要事务,至少 3 个 broker(事务状态 topic 副本因子默认 3)

# 消费者侧(下游应用):不设这个,EOS 形同虚设
isolation.level=read_committed
commit.interval.msEOS 把它从 30000ms 砍到 100ms——提交越频繁、未提交批越小、端到端延迟越低,但事务开销和吞吐损失随之上升;调它就是在延迟和吞吐间权衡。 isolation.level下游不设 read_committed 就会读到未提交甚至已回滚的记录,上游的 EOS 白做。
表 2.6 · 精确一次怎么做:at-least-once + 消费侧幂等 vs EOS v2
方案优势为什么没选 / 选中
at-least-once + 下游按业务键幂等去重 无事务开销、吞吐高;不依赖 ≥3 broker 每个下游都得自己实现去重(去重表 / 幂等 upsert),漏一处就重复;对"状态聚合"这类内部计算无能为力——状态会被重放放大
eos-v1(exactly_once,一 task 一 producer) 同样的端到端原子语义 已淘汰:高分区数下 producer 数量随 task 爆炸,内存与事务协调开销失控
exactly_once_v2(一线程一 producer,KIP-447) offset+状态+输出 原子提交,状态聚合也精确一次;producer 数只随线程数增长 选中:3.0 起默认;代价是提交延迟(commit.interval.ms 默认 100ms)+ 仅限单 Kafka 集群内、不覆盖外部副作用

带来的代价

① 延迟与吞吐:EOS 把 commit.interval.ms 默认从 30000ms 改到 100ms——这才是"开了 EOS 就变慢"的真因(更频繁提交 + read_committed 消费者 15–30% 吞吐下降),而非事务本身有多重。调大它换吞吐、调小换延迟(04 章)。② 边界止于 Kafka:事务只在单个 Kafka 集群内原子,不覆盖外部副作用——往 DB / HTTP / 发短信,事务回滚也撤不回,仍会重复。要端到端一致得上 outbox 模式(把外部写也变成 Kafka 写,由另一个消费者按幂等落地)。③ 部署门槛:需要至少 3 个 broker(事务状态 topic 副本因子默认 3),下游必须配 read_committed 否则 EOS 形同虚设。这与主教程的投递语义一脉相承——Streams 的 EOS 是消费组 + 事务在流处理上的封装,不是新机制。

想一想

一个 EOS v2 应用对每条记录的副作用是:写一条结果到输出 topic + 调一次外部支付 HTTP 接口。事务回滚(abort)后,输出 topic 那条记录消失了,但支付接口已经被调过。这条记录会怎样?problem 在哪、怎么修?

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

输出 topic 的写入随事务回滚被撤销,下游 read_committed 消费者看不到它——这部分精确一次成立。但外部支付 HTTP 调用不在 Kafka 事务里,已经发生且无法回滚;事务重试时这条记录会被重新处理、支付接口被再调一次,造成重复扣款。

根因:EOS 的原子边界止于单个 Kafka 集群,外部副作用在边界之外。修法是 outbox 模式——处理时不直接调支付,而是把"待支付"事件写进输出 topic(这步在事务内,精确一次),由一个独立消费者读出后调支付接口,并对支付接口做幂等(按 event-id 去重),把"恰好一次的外部效果"交给幂等而非事务来保证。

2.7串起来:一条记录在有状态聚合里的完整一生

六个机制不是孤立的。把它们串在一个具体场景上看协同:一条订单事件进入一个 EOS v2 的有状态聚合拓扑(按用户 ID 累计消费额),从进入到被下游读到,依次穿过状态、重分区、提交、投递四道关。

洞察 · 一条记录的旅程

① 进入与重分区(§2.3):订单事件以订单号为 key 进入;聚合要按用户 ID,于是 groupBy(userId) 改了 key,下游的 aggregate 触发一次 repartition——记录被 produce 进 repartition topic、再 re-consume 回来,此时相同用户 ID 已落同一 task。② 更新 RocksDB + 写 changelog(§2.1):聚合算子把该用户的累计额更新进本地 RocksDB,同一步把这次变更追加进 changelog topic。③ EOS 事务提交(§2.6):到 commit.interval.ms(EOS 下 100ms)时,这条记录的输入 offset 提交 + changelog 写 + 输出 produce 被绑进同一个 Kafka 事务一起提交——三者同生共死。④ 下游 read_committed 读(§2.6):下游消费者只读已提交事务里的输出,跳过任何被回滚的部分,拿到精确一次的累计结果。⑤ 若此刻实例崩溃(§2.1 + §2.2):再平衡把这个有状态 task 搬到别的实例,新 owner 回放 changelog(或从 standby 追尾)重建该用户的累计额——因为 RocksDB 只是缓存、changelog 才是真相,状态精确恢复到最后一次提交的位置,不重不漏。

这条链解释了一个常被误解的现象:为什么"开了 EOS 的有状态聚合"在崩溃恢复后既不丢数也不重复算。不是因为某个魔法标志,而是因为状态的真相在 changelog、changelog 的写又被事务和 offset 提交绑在一起——恢复时回放到的状态,和已提交给下游的输出,永远对齐在同一个事务边界上。

§本章 self-check

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

  1. 为什么说本地 RocksDB 只是缓存、changelog 才是真相源?故障转移时新 owner 从哪里恢复状态,standby 副本把什么从"分钟"压成了"秒"?
  2. selectKey 被称为 lazy,是什么意思?什么时候它身后才真的冒出一条 repartition topic?为什么优先用 groupByKey 而非 groupBy?
  3. stream-time 由什么驱动?为什么一个空闲分区会让窗口"永远不发结果"?这是 event-time 换来了什么而付的代价?
  4. 四种 join 里哪些要求共分区、哪些不要求?"分区数对得上但 join 零输出且无报错"嫌疑最大的原因是什么?
  5. EOS v2 的事务到底包住了哪三件事?为什么说"开了 EOS 变慢"的真因是 commit.interval.ms 而不是事务本身?EOS 的原子边界到哪里为止?
  6. (跨机制综合)一条记录进入"EOS v2 + 按用户 ID 聚合"的拓扑,从进入到被下游 read_committed 读到,依次穿过哪些机制?如果在事务提交后、下游读到前实例崩溃,状态如何保证不重不漏地恢复?
答案(先做完再展开)
  1. 本地 RocksDB 随实例消失(重调度/坏盘/再平衡搬走就没了),而每次 put 都同步追加进 compacted changelog——后者跨实例持久,能从零重建 RocksDB,所以它才是真相源。故障转移时新 owner 回放 changelog 重建状态(不是从旧实例的盘恢复,盘已不在)。standby 副本在另一实例后台尾随同一条 changelog 保温本地 RocksDB,接管时只需追上最后一小段增量,把全量回放(分钟)压成追尾(秒)。
  2. lazy 指 selectKey 只置"需重分区"标志、自己不立即重洗也不建 topic。只有当下游出现真正依赖 key 共分区的 stateful 算子(aggregate/count/join)时,Streams 才在它前面物化出一条 repartition topic。groupByKey 不改 key、不触发重分区;groupBy 改 key、要付一次 produce→re-consume 往返——能用前者就别用后者。
  3. stream-time = 该 task 见过的最大记录时间戳,单调不减,只在有记录到达时前进。空闲分区没有新记录,stream-time 冻结,依赖它推进到 windowEnd+grace 才关闭的窗口永远不关、结果永远不发。这是 event-time 换来可重放确定性(同输入同结果、与运行时刻无关、正确处理乱序回填)所付的代价——Streams 不看墙上时钟。
  4. KStream-KStream / KStream-KTable / KTable-KTable 等值 join 要求共分区(相同分区数 + 相同分区方式);GlobalKTable join 和外键 join 不要求(前者每实例全量副本,后者用隐藏 topic 重路由)。"分区数对得上但零输出无报错"嫌疑最大的是分区方式不一致——Streams 启动只校验分区数量、不校验分区器,相同 key 落到不同分区号,匹配全部静默丢失。
  5. 包住三件事:输入 offset 提交 + 状态 changelog 写 + 输出 produce,绑进一个 Kafka 事务原子提交/回滚,下游 read_committed 跳过回滚部分。"变慢"真因是 EOS 把 commit.interval.ms 默认从 30000ms 改到 100ms(提交更频繁 + read_committed 吞吐降),不是事务固定开销。原子边界止于单个 Kafka 集群,不覆盖 DB/HTTP/短信等外部副作用(要 outbox + 幂等)。
  6. 依次穿过:repartition(groupBy(userId) 改 key→repartition topic 重洗,§2.3)→ 状态写(更新本地 RocksDB + 追加 changelog,§2.1)→ EOS 事务提交(offset+changelog+输出 三者绑一个事务,§2.6)→ 下游 read_committed 读(只读已提交,§2.6)。提交后、下游读到前崩溃:再平衡把有状态 task 搬到新实例(§2.2),新 owner 回放 changelog 或从 standby 追尾重建该用户累计额(§2.1);因为状态真相在 changelog、changelog 写又与 offset 提交绑在同一事务,恢复到的状态与已提交给下游的输出对齐在同一事务边界,不重不漏。
进阶挑战 · 刚好够不着

给一个"窗口聚合 + EOS + 维表富集"的拓扑画出全部内部 topic

设计这样一条拓扑:从 orders(按订单号分区)读流,selectKey 改成用户 ID,做一个 5 分钟 tumbling 窗口的金额求和(带 grace + suppress 只发最终结果),再用一个小的 users 维表富集出用户等级,整条拓扑开 exactly_once_v2,结果写到 user-spending。不写代码,只回答:这条拓扑会让 Streams 在 broker 上建出哪些内部 topic(repartition / changelog 各几条、分别服务谁)?维表用 KTable 还是 GlobalKTable 才能免去对 orders 的共分区要求?suppress 的状态存活在哪、它的 changelog retention 要 ≥ 多少?

提示(卡住再展开)

顺着记录流向数:selectKey(userId) 后的窗口聚合是 stateful 且依赖 key→触发一条 repartition topic;窗口聚合的状态存储 + suppress 的缓冲各自有 changelog topic。维表若用 KTable 做 KStream-KTable join,要求与(已按 userId 重分区的)流共分区;用 GlobalKTable 则每实例全量、免共分区——但要权衡 users 是否够小。suppress 的 changelog retention 要覆盖"窗口大小 + grace"否则窗口没关状态先被清。把每条内部 topic 的命名前缀 <application.id>-… 写出来,再用 topology.describe() 的输出对照验证。