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 章对应机制——记不清机制本身时点回去补,本章不重讲机制,只讲机制被违反后会观察到什么。

严重度 ↑ 类别 → 恢复/再平衡 正确性 资源/语义 高 · 拖垮可用性 中 · 结果错或变慢 低 · 配置/边界 ① 恢复卡顿 回放 changelog ⑦ 扩不动 上限=分区数 ④ join 无输出 未共分区 ⑤ serde 崩溃 默认 serde ③ 隐藏重分区 改 key ⑥ OOMKilled 堆外 RocksDB ⑧ EOS 慢/重复 100ms+外部 ⑥' 迟到丢失 过 grace 改 key→延迟+存储
图 4.0八个失败模式按类别(横轴)× 严重度(纵轴)落位;编号与下文 §4.1–§4.8 一一对应。注意:右上角朱红的 ⑥ OOMKilled 是最反直觉的一条——它落在"高严重度"是因为容器被内核杀掉、堆指标却全程正常;底部那条朱红曲线提醒"改 key(③)"是延迟与存储成本的同一个源头。

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 头号运维痛点。容器编排里"重启即丢盘"会让每次重启都触发全量回放。

修复

streams.properties Properties
# 错误:默认 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 指标)再滚下一个实例。

陷阱 · stream-time 不是唯一会"冻结"的东西

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 则多出的线程闲置。并行天花板就是分区数,加再多实例也越不过去。

修复

scaling.properties Properties
# 现状: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 讲的正是这个自动机制及其隐藏成本)。默认正确,但成本不可见。

修复

RepartitionFix.java Java
// 错误: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 是 lazy 的,别被它骗了

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 与共分区 讲的就是这条静默失败路径)。

修复

JoinFix.java Java
// 错误: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 漂移让历史消息反序列化失败。前两条都源自"让运行时替你猜类型/猜归属"。

修复

SerdeNullFix.java Java
// 错误:依赖全局默认 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 的视野。

容器内存上限(内核盯着这条线) JVM 进程 堆 (heap) -Xmx 管这里 GC 看着正常 used heap 离 -Xmx 还远 空余(误导信号) 堆外 (off-heap) RocksDB · 每 store 一份 -Xmx 管不到 block cache ≈50MiB memtable 写缓冲 index/filter 块 × N 个 store 顶破上限 堆+堆外 > 容器 → 内核 OOMKill (exit 137) 调大 -Xmx → 堆吃更多额度 → 堆外更挤 OOMKill 反而更早
图 4.1OOMKill 的真因在于堆外那一列:RocksDB 的 block cache / memtable / index 块按"每 store 一份 × N 个 store"在 JVM 堆之外累加,-Xmx 完全看不到它。注意:左侧"堆还有空余"是个诱饵信号——内核盯的是朱红虚线那条"堆+堆外"总和的容器上限,所以调大 -Xmx 只会挤掉堆外额度、让 OOMKill 更早到。

修复

BoundedRocksDBConfig.java Java
// 错误:不管 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) {}
}
deploy.env / streams.properties Properties
# 挂上 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 才反映堆 + 堆外的真实占用。

陷阱 · JVM 堆是 OOMKill 的诱饵

容器被 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"最危险。

修复

GraceFix.java Java
// 错误:没显式声明 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 内精确一次"。

修复

EosFix.java / streams.properties Java
// 配置: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。

陷阱 · EOS 的"精确一次"只在 Kafka 边界内成立

exactly_once_v2 绑定的是 offset 提交 + changelog 写 + 输出 produce 三者的 Kafka 内原子性,下游必须 read_committed 才看得到这层保证。它不覆盖任何外部系统:processor 里的 DB 写、HTTP 调用、发短信,在事务重试时都会重复。面试里把 EOS 说成"端到端精确一次"是减分项——正确说法是"Kafka 内事务性精确一次,外部要靠幂等/outbox"。

4.D反模式与"何时不该用 Streams"

八个失败模式之外,还有几条不在单条陷阱里、但会反复出现的反模式,以及一个更诚实的问题:什么时候根本不该上 Streams。

表 4.1 · 反模式速查
反模式为什么是错的正确做法
用 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

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

  1. 容器报 exit code 137 反复重启,但 JVM used heap 离 -Xmx 还很远、GC 平静。最可能的根因是什么?为什么"调大 -Xmx"会让情况更糟?应该看哪个指标判断?
  2. 两个流 join,逻辑正确、进程健康、零异常,输出 topic 却一条都没有。第一个该怀疑的是什么?分区"数"不符和分区"方式"不符,哪一种更难查、为什么?
  3. 同事说"我只在代码里 selectKey 了一次,怎么会有 repartition topic 的存储和延迟成本?" 他的归因错在哪?用什么命令能数清实际建了几个内部 topic?
  4. (设计层)某团队把 processing.guarantee 切到 exactly_once_v2,期望"用户绝不会被重复扣款"。这个期望哪里站不住?要真正做到"不重复扣款",正确的架构是什么?
答案(先做完再展开)
  1. 根因是 RocksDB 在堆外分配内存(block cache + memtable + index/filter,每 store 一份),-Xmx 管不到它;容器上限盖的是"堆 + 堆外"总和,堆还有空余但总和已顶破上限,被内核 OOMKill。调大 -Xmx 把更多额度划给堆、留给堆外的更少,OOMKill 反而更早。判断应看容器 RSS 对容器 limit,不是 used heap 对 -Xmx。修复:RocksDBConfigSetter 共享 LRUCache + WriteBufferManager 把堆外钉成常数(§4.A 陷阱 ⑥、图 4.1)。
  2. 第一个该怀疑共分区:等值 join 要求两侧相同分区数 + 相同分区方式。分区数不符会让拓扑启动直接报错(fail-fast,反而好查);分区方式不符(数相同、分区器不同)校验通过却把记录路由到不同 task,匹配全丢、零输出零报错——静默失败更难查。修复:join 前对一侧显式 repartition 对齐,或小维表用 GlobalKTable(陷阱 ④)。
  3. 归因错在:selectKey 是 lazy 的,单独存在不建任何 topic,只打"key 变了"的标记;真正物化 repartition topic 的是下游依赖 key 的算子(groupByKey / join / aggregate)。所以数源码里 selectKey 的次数没意义,要用 topology.describe() 看实际生成了几个内部 topic(陷阱 ③)。
  4. 期望站不住,因为 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 分区是并行上限,线程/实例别超。把这些连起来就是一份体检报告。