Chapter 04
陷阱与失败模式
03 章把 DSL 跑通了,前三章合起来交付了一个能用的 Streams 应用:消费组撑起 task、本地 RocksDB 存状态、changelog 当真相源、窗口靠 stream-time 推进、EOS 把读-处理-写包进一个事务。那是机制按设计工作的样子。本章是同一批机制被违反时的样子——八个具名失败模式,每个都能回链到 02 章某条原理。读完应能在症状出现前就认出隐患,而不是等到生产环境 OOMKilled 才回头查内存模型。
本章你将建立的 schema
- 每个生产事故都是某条 02 章原理被违反的可观测后果——症状是表象,根因在机制
- 三类失败模式:状态恢复/再平衡、正确性、资源/语义——对应三组锚点
- 最反直觉的两条:OOMKilled 时 JVM 堆是诱饵(真因在堆外 RocksDB),join 无输出时没有任何报错(共分区静默失败)
- 修复几乎都是"显式声明本来被自动推断的东西":显式 Serde、显式 grace、显式 LRUCache、显式 standby
八个失败模式按归因分三组。每组先给一张地图(见 §4.0 分类图),再逐条拆 症状 → 根因 → 修复 → 如何避免。根因一律回链 02 章对应机制——记不清机制本身时点回去补,本章不重讲机制,只讲机制被违反后会观察到什么。
4.A状态恢复与再平衡类
这组失败模式的共同根:有状态 task 的状态住在本地 RocksDB,真相源却是 Kafka 里的 changelog 日志(回 §2.1 状态存储)。task 一旦换了 owner,新 owner 手上没有本地副本,只能从头回放 changelog 重建——状态越大、带宽越窄,这段重建越长。缩放也卡在同一处:并行度被分区数钉死(回 §2.2 再平衡与缩放)。
陷阱 ① · 再平衡时状态恢复卡顿数分钟
症状:滚动重启或扩容后,新实例的 task 卡在 RESTORING 状态数分钟,期间这些分区的输入完全不处理;端到端延迟出现分钟级尖峰;日志里刷 Restoration in progress,恢复速率受限于 changelog 读取带宽。重启越频繁,停顿越频繁。
根因:有状态 task 被重新分配给新实例时,新 owner 本地没有 RocksDB 副本,必须从 changelog topic 的最早 offset 回放、逐条重建状态存储后才能开始处理(§2.1 讲的就是这条恢复路径)。多 GB 状态 + 受限带宽 = 分钟级阻塞,这是 Streams 头号运维痛点。容器编排里"重启即丢盘"会让每次重启都触发全量回放。
修复
# 错误:默认 0 个 standby,任何故障转移都从零回放 changelog
# num.standby.replicas=0
# 正确:至少 1 个热备,故障转移从"分钟回放"降到"秒级追尾"
num.standby.replicas=1
# rack-aware standby(KIP-925,3.2+):把 standby 摆到另一可用区,
# 同时挡住"主实例挂 + 同机架 standby 一起挂"
rack.aware.assignment.tags=zone
client.tag.zone=az-a
# 提高恢复并行度:恢复用独立线程拉 changelog,不占处理线程
# (默认即开,确认未被调低)
另外两条在部署层,不在配置文件:
- 持久卷(StatefulSet + PVC):让 RocksDB 目录跨 Pod 重启存活。重启后本地状态还在,Streams 只补 changelog 尾部增量,而非全量回放——这是把"每次重启分钟级"压回"秒级"的最大杠杆。
- static membership(
group.instance.id):给实例固定身份,滚动重启时 group coordinator 不把它当成新成员、不触发 task 重分配。但注意见下方警告——KIP-1071 新再平衡协议下它还不支持。
如何避免:上线前就按"状态大小 ÷ changelog 读带宽"估算最坏恢复时长,写进容量规划;把 num.standby.replicas≥1 和持久卷设成有状态 Streams 应用的默认基线,而不是出事后才补。把"无状态滚动重启"和"有状态滚动重启"当成两种运维动作——后者必须先确认 standby 已追平(看 standby-tasks 指标)再滚下一个实例。
static membership 能消掉滚动重启的再平衡,但 KIP-1071 的新 Streams 再平衡协议(4.2 起新集群默认开)当前还不支持 static membership、也不支持运行中改 topology、不支持在线 classic→streams 迁移(截至 2026-06,Kafka 4.3)。如果新集群默认开了 streams.version=1,别假设旧的 static membership 配置还生效——先确认协议版本,再决定靠 standby 还是靠 static membership 压恢复时间。
陷阱 ⑦ · 加机器也扩不过 N 个消费者
症状:输入 topic 有 6 个分区,部署到 10 个实例(或把 num.stream.threads 开到很大),却发现只有 6 份工作在跑,多出来的实例/线程 ASSIGNED 了零个 task、CPU 全程空转;吞吐怎么加机器都不涨。
根因:Streams 没有独立调度器,并行单位是 task,而 task 数 = 子拓扑最大输入分区数,运行时不可改(§2.2)。task 多于线程则一个线程多路复用几个 task;线程多于 task 则多出的线程闲置。并行天花板就是分区数,加再多实例也越不过去。
修复
# 现状:6 分区输入 → 最多 6 个 task
# 错误:单实例把线程开到 16,超过 6 的 10 个线程全闲置
# num.stream.threads=16
# 正确:先把每实例线程提到接近 task 数,再横向加实例直到 task 摊满
# 例如 2 实例 × 3 线程 = 6 个 slot,正好 1 task/slot
num.stream.threads=3
# 然后部署 2 个实例
# 要超过 6 路并行:只能重分区输入 topic(破坏性,需停机/双写迁移)
# kafka-topics --alter --partitions 12 ... ← 不可逆地改变 key→分区映射
如何避免:在创建输入 topic 时就按未来吞吐峰值定分区数——分区是缩放上限,事后加分区会打乱 key→分区 映射、破坏既有有状态算子的共分区前提,代价远高于一开始多给几个分区。监控时看 assigned-tasks 而非实例数:实例数涨、assigned-tasks 不涨,就是撞到了分区天花板。
4.B正确性类
这组的危险在于不报错:拓扑照常运行、指标一片绿,但结果悄悄偏了——多了隐藏 topic、少了本该有的 join 输出、或在某条 null key 上崩在一个看似无关的位置。根都在"改 key"与"共分区"两条机制(§2.3 repartition / §2.5 join)。
陷阱 ③ · 隐藏的 repartition topic 抬高存储与延迟
症状:拓扑里没显式建任何 topic,broker 上却冒出一批 <app-id>-<name>-repartition 内部 topic,吃掉额外磁盘;端到端延迟比预期高一截,因为每条记录在聚合/join 前多走了一次 produce→re-consume 往返。改 key 的算子越多,往返越多。
根因:聚合和等值 join 要求相同 key 落在同一 task/分区。任何改 key 的操作(selectKey / map / flatMap / groupBy)都会给下游打上"需重分区"标志,于是下游 stateful 算子自动建一个 repartition topic 把数据重洗一遍(§2.3 讲的正是这个自动机制及其隐藏成本)。默认正确,但成本不可见。
修复
// 错误:map 改了 key(哪怕 value 没动),再 groupBy 又改一次 →
// 触发一次甚至两次重分区往返
KStream<String, Order> s = orders
.map((k, v) -> KeyValue.pair(v.userId(), v)) // 改 key → 标记需重分区
.groupBy((k, v) -> v.userId()) // 又改 key → 又一次重分区
.count(...); // 实际建了多余的 repartition topic
// 正确①:value 变换用 mapValues(不碰 key,不触发重分区)
KStream<String, Long> amounts = orders
.mapValues(Order::amount); // key 不变,零重分区
// 正确②:key 已经对了就用 groupByKey(绝不重分区),别用 groupBy
KTable<String, Long> perUser = orders
.selectKey((k, v) -> v.userId()) // selectKey 是 lazy:自己不建 topic
.groupByKey(Grouped.with(Serdes.String(), orderSerde))
.count(); // 这里才物化一次重分区,且仅一次
// 正确③:必须重 key 时,用 Repartitioned 命名 + 控分区数,让它可观测可控
KStream<String, Order> rk = orders
.selectKey((k, v) -> v.userId())
.repartition(Repartitioned.<String, Order>as("orders-by-user")
.withNumberOfPartitions(12));
如何避免:把 topology.describe() 的输出当成 code review 的一部分——它会列出实际生成的每一个 repartition / changelog 内部 topic,数量对不上预期就说明有多余的改 key 操作。习惯性优先 mapValues 而非 map、groupByKey 而非 groupBy,只在 key 真的需要变时才 selectKey。
selectKey 单独存在时不建任何 topic——它只是给流打上"key 变了"的标记。真正物化 repartition topic 的是下游依赖 key 的算子(groupByKey / join / aggregate)。所以"我只 selectKey 了一下,怎么会有 repartition 成本"是错的归因:成本由下游触发,看 topology.describe() 才数得准,数源码里 selectKey 的次数没用。
陷阱 ④ · join 静默无输出
症状:两个流/表 join,逻辑看着完全正确,输出 topic 却一条都没有,且没有任何异常或错误日志。改 key、改 join 窗口都不见效。最折磨人的是它不崩——只是安静地什么都不发。
根因:等值 join 要求两侧共分区——相同分区数 + 相同分区方式(同一个 producer 分区器把同一 key 落到同一编号分区)。分区数不一致时拓扑启动会直接失败;但分区方式不一致(分区数相同、分区器不同)时,Streams 校验通过、却把本该匹配的记录路由到不同 task,匹配全部丢失,零输出零报错(§2.5 join 与共分区 讲的就是这条静默失败路径)。
修复
// 错误:clicks 上游被某个外部 producer 用了不同分区器写入,
// 与 impressions 分区数相同但分区方式不同 → join 静默丢匹配
KStream<String, Joined> bad = impressions
.join(clicks,
(imp, clk) -> merge(imp, clk),
JoinWindows.ofTimeDifferenceAndGrace(
Duration.ofMinutes(5), Duration.ofMinutes(1)));
// 正确①:join 前对至少一侧显式 repartition,强制 Streams 用自己的
// 默认分区器重洗,对齐分区方式(也对齐分区数)
KStream<String, Click> alignedClicks = clicks
.repartition(Repartitioned.<String, Click>as("clicks-aligned")
.withNumberOfPartitions(impressionsPartitions));
KStream<String, Joined> good = impressions.join(alignedClicks, ...);
// 正确②:维表小且慢变,用 GlobalKTable —— 每实例全分区副本,
// 天然免共分区,是这条陷阱的逃生口
GlobalKTable<String, User> users = builder.globalTable("users");
KStream<String, Enriched> enriched = events.join(
users,
(eventKey, event) -> event.userId(), // 从流记录提 join key
(event, user) -> enrich(event, user));
如何避免:把"join 前两侧共分区"列成上线 checklist 的硬项——确认两侧分区数相同,且要么都由本 Streams 应用写入(同一默认分区器),要么对其中一侧显式 repartition 重洗。小维表一律考虑 GlobalKTable 或外键 join(KIP-213),两者都免共分区。监控 join 输出 topic 的 records-out:长期为 0 而输入不为 0,第一个怀疑对象就是共分区。
分区数不符会让拓扑启动报错——这种 fail-fast 反而友好。真正吃人的是分区方式不符:校验通过、进程健康、指标全绿,唯独结果是空的。遇到"join 没输出又不报错",别先怀疑业务逻辑或窗口大小,先验共分区。
陷阱 ⑤ · serde 与 null key 崩溃
症状:运行中抛 ClassCastException 或 SerializationException,栈顶指向某个序列化器,但报错位置和你"觉得"该出错的算子对不上;或者聚合结果莫名少了一批记录,怎么查都查不到它们去哪了。schema 升级后老消息开始反序列化失败。
根因:三件事。其一,依赖 default.key/value.serde 但某个分支的类型和默认 serde 不符——Streams 会用默认 serde 去序列化它,类型不匹配就在那个算子炸。其二,key 为 null 的记录在进入聚合/join 前会被静默丢弃(有状态算子按 key 分组,null key 无处归属),表现为"结果少了一批"而非报错。其三,schema 漂移让历史消息反序列化失败。前两条都源自"让运行时替你猜类型/猜归属"。
修复
// 错误:依赖全局默认 serde;不同分支类型不同时会用错 serde 炸;
// null key 直接进 groupByKey 被悄悄丢
KStream<String, Payment> payments = builder.stream("payments"); // 用 default serde
KTable<String, Long> bad = payments
.groupByKey() // null key 记录在此静默消失
.count();
// 正确:每个 source/sink/store 显式传 Serde;先过滤 null key 再聚合
KStream<String, Payment> typed = builder.stream(
"payments",
Consumed.with(Serdes.String(), paymentSerde)); // 显式 key/value serde
KTable<String, Long> good = typed
.filter((k, v) -> k != null && v != null) // 显式处理 null,而非让它静默蒸发
.groupByKey(Grouped.with(Serdes.String(), paymentSerde))
.count(Materialized.with(Serdes.String(), Serdes.Long())); // store 也显式 serde
如何避免:把"显式 Consumed / Produced / Grouped / Materialized 传 Serde"当成纪律,少依赖全局默认——默认 serde 是"省事但把类型错误推迟到运行时"的交易。在聚合/join 前显式决定 null key 该丢还是该补默认值,让它成为代码里看得见的一行。用 Schema Registry 的兼容模式(backward/forward)兜住 schema 演进。
4.C资源与语义类
最后一组踩的是"机制在 JVM 之外"和"语义边界在你以为之外":RocksDB 的内存不归 -Xmx 管,stream-time 只在有数据时前进,EOS 的原子性只覆盖 Kafka 之内。三条都因为读者把边界画错了地方。
陷阱 ⑥ · 容器 OOMKilled,但 JVM 堆指标全程正常
这条最反直觉,讲透。症状:容器被内核 OOMKill(exit code 137),Pod 反复重启;可你盯着 JVM 堆指标——used heap 离 -Xmx 还很远,GC 日志平静,full GC 没几次。把 -Xmx 调大反而让 OOMKill 来得更快。监控面板上 JVM 一切正常,容器却在被杀。
根因:RocksDB 在堆外(off-heap)分配内存,-Xmx 完全管不到它。每个有状态 task 的每个 store 都有自己的 block cache(默认约 50 MiB)+ memtable(写缓冲)+ index/filter 块,全在 native 内存。一个应用几十个有状态 task,堆外内存能轻松到几个 GB。容器内存上限同时盖住"堆 + 堆外 + 线程栈 + 元空间";JVM 堆离 -Xmx 还远,但堆 + 堆外的总和已经顶破容器上限,内核就 OOMKill 整个进程。把 -Xmx 调大等于把更多额度划给堆、留给堆外的更少,OOMKill 反而更早——这就是"调大堆反而更糟"的来由。根在 §2.1:选 RocksDB 是为本地亚毫秒读,代价之一就是这部分内存逃出了 JVM 的视野。
-Xmx 完全看不到它。注意:左侧"堆还有空余"是个诱饵信号——内核盯的是朱红虚线那条"堆+堆外"总和的容器上限,所以调大 -Xmx 只会挤掉堆外额度、让 OOMKill 更早到。修复
// 错误:不管 RocksDB 内存,每个 store 各自一份 block cache + memtable,
// 堆外随 task 数线性膨胀,最终顶破容器上限被 OOMKill
// 正确:自定义 RocksDBConfigSetter,让所有 store 共享一个有上限的
// LRUCache + WriteBufferManager —— 把堆外内存钉成一个常数
public class BoundedRocksDBConfig implements RocksDBConfigSetter {
// static:进程内所有 task 的所有 store 共用这两个对象
private static final long TOTAL_OFF_HEAP = 256L * 1024 * 1024; // 256 MiB 总上限
private static final long TOTAL_MEMTABLE = 128L * 1024 * 1024;
private static final Cache CACHE = new LRUCache(TOTAL_OFF_HEAP);
private static final WriteBufferManager WBM =
new WriteBufferManager(TOTAL_MEMTABLE, CACHE); // memtable 也计入同一池
@Override
public void setConfig(String storeName, Options options,
Map<String, Object> configs) {
BlockBasedTableConfig table =
(BlockBasedTableConfig) options.tableFormatConfig();
table.setBlockCache(CACHE); // 共享 cache
table.setCacheIndexAndFilterBlocks(true); // index/filter 也进 cache,纳入上限
options.setWriteBufferManager(WBM); // 共享写缓冲池
options.setTableFormatConfig(table);
}
@Override public void close(String storeName, Options options) {}
}
# 挂上 config setter
rocksdb.config.setter=com.example.BoundedRocksDBConfig
# 容器内存 = 堆(-Xmx) + 堆外上限(LRUCache+WBM) + 线程栈/元空间 + 余量
# 例:-Xmx2g + 256MiB cache + ~512MiB 余量 → 容器 limit ≥ 2.75g
# 抑制 glibc 多 arena 造成的 native 内存碎片(容器里尤其明显)
# 作为环境变量设置:MALLOC_ARENA_MAX=2
如何避免:把容器内存预算显式拆成"堆 + 堆外缓存上限 + 余量"三项分别算,而不是"给个 -Xmx 再加点"。任何上 RocksDB 的有状态 Streams 应用,默认就配 RocksDBConfigSetter 共享 LRUCache + WriteBufferManager,把堆外钉成常数——否则堆外随 task 数线性涨,迟早撞上限。监控容器 RSS(而非只看 JVM 堆),RSS 才反映堆 + 堆外的真实占用。
容器被 OOMKill 时第一反应往往是"加 -Xmx"或"看 GC"——两条都走偏。RocksDB 的内存在堆外,-Xmx 与 GC 指标对它一无所知。判断依据应是容器 RSS 对容器 limit,不是 used heap 对 -Xmx。exit code 137 + 堆指标正常 = 几乎必是堆外(RocksDB 或 直接 ByteBuffer)顶破了上限。
陷阱 ⑥' · 迟到事件无声消失
症状:明明收到了带较早事件时间的记录,窗口聚合结果却没把它算进去,且没有任何错误——只在 dropped-records 指标上加了 1(不看这个指标就完全无感)。补数据、回放历史时尤其常见:旧时间戳的记录大批被丢。
根因:窗口有宽限期 grace,一旦 stream-time ≥ window_end + grace,迟到记录被静默丢弃,只计入 dropped 指标(§2.4 窗口与时间)。再叠加两个放大因素:event-time 与 processing-time 混淆导致时间判断错位;以及历史上某些默认 grace 行为会无声丢弃很久——所以"没显式设 grace"最危险。
修复
// 错误:没显式声明 grace,迟到记录何时被丢由默认行为决定,不可控
TimeWindows badWin = TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5));
// 正确:显式声明窗口大小 + grace,把"容忍迟到多久"写成代码里看得见的决策
TimeWindows win = TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(5), // 窗口大小
Duration.ofMinutes(10)); // grace:窗口关后还容忍 10 分钟的迟到
KTable<Windowed<String>, Long> counts = events
.groupByKey()
.windowedBy(win)
.count()
// suppress:只在窗口彻底关闭后发一个最终结果,避免下游看到中间态
.suppress(Suppressed.untilWindowCloses(
Suppressed.BufferConfig.unbounded()));
// 同时:changelog/窗口 store 的 retention 必须 ≥ 窗口 + grace,否则
// 状态先于 grace 被清掉,迟到记录连匹配的窗口都找不到
// Materialized.as(...).withRetention(Duration.ofMinutes(20))
如何避免:永远用 ofSizeAndGrace 显式声明 grace,禁止依赖默认;把 grace 当成一条业务决策("业务能容忍多久的迟到")而非一个被忽略的参数。给 dropped-records 指标配告警——它从 0 开始涨,就是有迟到事件在被丢。窗口 store 的 retention 设成 ≥ 窗口 + grace,并在回放历史/补数时单独评估 grace 是否够大。
陷阱 ⑧ · 开 EOS 后延迟变高或外部副作用意外重复
症状:把 processing.guarantee 切到 exactly_once_v2 后,端到端延迟明显上升、下游 read_committed 消费者吞吐降 15-30%;或者更糟——你以为"精确一次"了,外部系统(数据库 / HTTP / 短信)却仍然收到重复调用。
根因:两件被画错的边界。其一,EOS 把 commit.interval.ms 默认从 30000ms 改成 100ms——频繁提交事务是延迟上升和吞吐下降的真因,不是"事务本身慢"(§2.6 EOS v2)。其二,EOS 的原子性只覆盖单个 Kafka 集群内的"offset 提交 + 状态 changelog 写 + 输出 produce",完全不延伸到外部副作用——你在 processor 里发的 HTTP / 写的 DB,在事务 abort 重试时会再执行一次。EOS 不是"全世界精确一次",只是"Kafka 内精确一次"。
修复
// 配置:EOS 下 commit.interval.ms 默认 100ms,按延迟/吞吐权衡调整
// (properties)
// processing.guarantee=exactly_once_v2
// commit.interval.ms=200 # 在延迟和吞吐之间找平衡,不必死守 100
// # EOS 需要 ≥3 个 broker(事务状态有副本要求)
// 错误:在 EOS 拓扑里直接做外部副作用,以为 EOS 会保证它精确一次
stream.foreach((k, order) -> {
paymentGateway.charge(order); // 事务 abort 重试时会重复扣款!EOS 管不到外部
});
// 正确:外部副作用走幂等或 outbox —— 让"重复"在外部侧被吸收
// ① 幂等:用 event-id 做外部 upsert,重复写无副作用
stream.foreach((k, order) ->
db.upsertByKey(order.eventId(), order)); // 同一 eventId 重复执行结果一致
// ② outbox:Streams 只写一条 Kafka 输出(受 EOS 保护),
// 由独立消费者带去重地投递给外部系统
KStream<String, Charge> outbox = orders
.mapValues(o -> new Charge(o.eventId(), o.amount()));
outbox.to("charge-outbox"); // 外部投递在另一进程按 eventId 去重
如何避免:开 EOS 前先做两道判断。第一,延迟预算能否吃下 100ms 量级的提交间隔——吃不下就调 commit.interval.ms 权衡,或重新评估是否真需要 EOS(很多场景 at_least_once + 下游幂等就够)。第二,拓扑里有没有外部副作用——有就一律走幂等(按 event-id upsert)或 outbox 模式,绝不指望 EOS 替你保证外部精确一次。部署时确认 broker ≥ 3。
exactly_once_v2 绑定的是 offset 提交 + changelog 写 + 输出 produce 三者的 Kafka 内原子性,下游必须 read_committed 才看得到这层保证。它不覆盖任何外部系统:processor 里的 DB 写、HTTP 调用、发短信,在事务重试时都会重复。面试里把 EOS 说成"端到端精确一次"是减分项——正确说法是"Kafka 内事务性精确一次,外部要靠幂等/outbox"。
4.D反模式与"何时不该用 Streams"
八个失败模式之外,还有几条不在单条陷阱里、但会反复出现的反模式,以及一个更诚实的问题:什么时候根本不该上 Streams。
| 反模式 | 为什么是错的 | 正确做法 |
|---|---|---|
| 用 GlobalKTable 装大表 | 每实例全量加载,磁盘+内存+启动 bootstrap 延迟都随表大小线性涨 | 只给"小而慢变"的维表用 GlobalKTable;大表用共分区的 KTable join 或外键 join |
| 在低流量/空闲分区上等窗口结果 | stream-time 只在有记录到达时前进;没数据则时间冻结、窗口永不关、结果不发 | 理解 stream-time 是数据驱动的;低流量场景考虑用心跳记录推进,或接受结果延迟到下一条记录 |
| 把 suppress 当"最终结果"保证,却用 emitEarlyWhenFull | 内存压力下 emitEarlyWhenFull 会发非最终的中间结果,违背"只发最终值"的预期 |
要严格最终结果用 shutDownWhenFull;用 emitEarlyWhenFull 就接受可能有中间态 |
| 运行中改 topology 还想热升级 | topology 变更改变内部 topic / task 划分;KIP-1071 新协议也不支持在线改 topology | topology 变更走停机部署 + 评估状态兼容;改 application.id 等于全新应用(状态从头建) |
| 无状态简单转发也上 Streams | 白白引入 changelog / 状态恢复 / 再平衡复杂度,换不来对应价值 | 纯过滤/转发/路由用原生 consumer-producer 就够;状态/窗口/join/EOS 才值得上 Streams |
什么时候不该用 Kafka Streams(诚实回答)
- 源和汇主要不是 Kafka:Streams 的容错、EOS、共分区全建在 Kafka 之上。如果数据主要来自数据库/对象存储、要写去多个异构外部系统,Streams 的优势用不上,反而要为它的 Kafka 中心假设买单。
- 需要跨多个 Kafka 集群的处理:EOS 只在单集群内原子(陷阱 ⑧),跨集群 join / 事务不在 Streams 的能力内。
- 复杂 CEP 或机器学习管线:需要复杂事件模式匹配、迭代计算、ML 推理编排时,专门的流处理引擎(Flink)或批/流统一框架更合适——它们有独立集群 + checkpoint 模型,承载得起这类负载。
- 只是简单消费:只要消费一个 topic 做点无状态处理,原生 consumer 更轻、心智负担更小(见上表最后一行)。
判别口径:要本地状态 / 窗口 / join / EOS,且源汇都在一个 Kafka 集群里——才是 Streams 的甜区。偏离这两条越远,越该换工具。这条判别会在 05 章综合实战里被反复用到。
§本章 self-check
先合上教程,把你能想到的答案写在纸上或编辑器里。 写完再点开答案对照——直接点开等于把这一节当再读一遍。
- 容器报
exit code 137反复重启,但 JVMused heap离-Xmx还很远、GC 平静。最可能的根因是什么?为什么"调大-Xmx"会让情况更糟?应该看哪个指标判断? - 两个流 join,逻辑正确、进程健康、零异常,输出 topic 却一条都没有。第一个该怀疑的是什么?分区"数"不符和分区"方式"不符,哪一种更难查、为什么?
- 同事说"我只在代码里
selectKey了一次,怎么会有 repartition topic 的存储和延迟成本?" 他的归因错在哪?用什么命令能数清实际建了几个内部 topic? - (设计层)某团队把
processing.guarantee切到exactly_once_v2,期望"用户绝不会被重复扣款"。这个期望哪里站不住?要真正做到"不重复扣款",正确的架构是什么?
答案(先做完再展开)
- 根因是 RocksDB 在堆外分配内存(block cache + memtable + index/filter,每 store 一份),
-Xmx管不到它;容器上限盖的是"堆 + 堆外"总和,堆还有空余但总和已顶破上限,被内核 OOMKill。调大-Xmx把更多额度划给堆、留给堆外的更少,OOMKill 反而更早。判断应看容器 RSS 对容器 limit,不是 used heap 对-Xmx。修复:RocksDBConfigSetter共享LRUCache+WriteBufferManager把堆外钉成常数(§4.A 陷阱 ⑥、图 4.1)。 - 第一个该怀疑共分区:等值 join 要求两侧相同分区数 + 相同分区方式。分区数不符会让拓扑启动直接报错(fail-fast,反而好查);分区方式不符(数相同、分区器不同)校验通过却把记录路由到不同 task,匹配全丢、零输出零报错——静默失败更难查。修复:join 前对一侧显式
repartition对齐,或小维表用GlobalKTable(陷阱 ④)。 - 归因错在:
selectKey是 lazy 的,单独存在不建任何 topic,只打"key 变了"的标记;真正物化 repartition topic 的是下游依赖 key 的算子(groupByKey/join/aggregate)。所以数源码里 selectKey 的次数没意义,要用topology.describe()看实际生成了几个内部 topic(陷阱 ③)。 - 期望站不住,因为
exactly_once_v2的原子性只覆盖单个 Kafka 集群内的 offset 提交 + changelog 写 + 输出 produce,不延伸到外部系统;processor 里直接调支付网关,在事务 abort 重试时会重复扣款。正确架构:外部副作用走幂等(按 event-id 对外部 upsert,重复无副作用)或 outbox(Streams 只写一条受 EOS 保护的 Kafka 输出,独立消费者按 event-id 去重后投递给支付系统)。同时注意 EOS 把commit.interval.ms默认改成 100ms 带来的延迟代价(陷阱 ⑧)。
给一个有状态 Streams 应用做"生产就绪体检"
设拓扑:orders(12 分区)经 selectKey(userId) 后与 users 表 join,再按 5 分钟滚动窗口聚合每用户下单额,开 exactly_once_v2 输出到下游并同时给风控系统发 HTTP 告警。部署在 Kubernetes,每 Pod -Xmx2g、容器 limit 2.5g、num.standby.replicas=0、users 用普通 KTable。请只凭本章的八个失败模式,列出这套配置里所有会出事的点,并给每个点一句话修复。目标是一次体检命中 ≥6 处。
提示(卡住再展开)
逐条对号入座:①恢复——standby=0 + 没提持久卷;③隐藏重分区——selectKey 下游 join/聚合会物化 repartition topic,topology.describe() 数一下;④join——users 与 orders 是否共分区(12 分区?同分区器?),不确定就 GlobalKTable 或 repartition 对齐;⑥OOMKill——容器 2.5g 只比 -Xmx2g 多 0.5g,没给 RocksDB 堆外留额度、也没 RocksDBConfigSetter;⑥'迟到——窗口有没有显式 grace + retention;⑧EOS——HTTP 告警是外部副作用,EOS 管不到、重试会重复发,且 100ms 提交间隔的延迟代价 + broker 是否 ≥3。⑦扩缩——12 分区是并行上限,线程/实例别超。把这些连起来就是一份体检报告。