Chapter 04

生产陷阱:机制被违反时的失败模式

02 章讲清了存储、分区顺序、副本/ISR、再平衡、投递语义这五套机制为什么这么设计。本章是同一批机制在被误用、被错配、或被默认值悄悄违反时的失败模式——每一条都能回溯到某个具体机制上的一次违约。

本章你将建立的 schema

  • 数据丢失、正确性、性能运维三类失败模式各自的根因机制不同——按根因排序而非按现象排序
  • "不丢消息"是 acks=all + min.insync.replicas + replication.factor + unclean.leader.election=false + 处理后再提交的组合命题,不是单个开关
  • 每条陷阱的三要素:症状(你观测到什么)→ 根因(违反了哪条机制)→ 修复(错误写法对正确写法)→ 如何在设计层面避免再次触发
  • 什么时候不该用 Kafka——失败模式的总和会反过来界定它的适用边界

这一章按根因与严重度排序,不按出现频率。数据丢失放最前:丢掉的记录不会回来,而重复和乱序还能在下游补救。每条失败模式都标注它违反了 02 章的哪条机制——失败模式不是孤立的故障,是某个设计取舍在被忽略时显形。

失败类别(按根因机制) 严重度 数据丢失 正确性 性能 / 运维 不可逆 可补救 可调优 acks=1 / unclean 已提交记录消失 重复 / 乱序 幂等可去重 再平衡风暴 / 倾斜 配置可缓解 因果链 ↘
图 4.0陷阱按"根因机制"(X 轴)与"能否补救"(严重度,Y 轴)二维分布;越靠左上越致命。 注意:数据丢失类落在"不可逆"带,所以本章把它放最前;右下虚线箭头表示一类失败会引发另一类——配错持久性会被误当成性能问题来调,根因找错。
洞察 · 按根因排序而非按现象

同一个现象("消息没了""消费者卡住")可能由多条机制失败造成。把陷阱挂在根因机制上,遇到现象时才能反查是哪条机制被违反,而不是逐个试配置。本章三节正对应三类根因:持久性机制(副本/ISR)、投递语义机制(幂等/顺序)、消费组协调机制(再平衡)。

4.1数据丢失类(最严重,不可逆)

数据丢失的共同根因是持久性保证被高估:producer 以为记录已经安全落地,实际上它只到了一台机器、或还没复制就因故障消失。全部回溯到 §2.3 副本与 ISR 里 HW(高水位)与 ISR 的语义。

陷阱 1:acks=1 丢数据(leader ack 后未复制即崩溃)

症状:producer 端 send() 全部返回成功、回调拿到了 offset;某次 broker 故障切换后,下游消费者读到的最高 offset 比 producer 确认过的要小——中间一段记录凭空消失,且无任何异常抛出。

根因:acks=1 表示 leader 把记录写进自己的日志就立即 ack,不等待 follower 复制。回链 §2.3:HW 只在 ISR 中所有副本都 fetch 到之后才推进。acks=1 在"leader 已写、follower 未同步"这个窗口里 ack 了 producer;若 leader 此刻崩溃,新 leader 从某个 follower 选出,那段没复制出去的记录随旧 leader 的磁盘一起丢失。这是用持久性换延迟。

producer-acks.properties Properties
# 错误写法:leader 单方面确认,复制前的崩溃窗口会丢已 ack 的记录
acks=1

# 正确写法:等所有 ISR 成员都复制到才确认
acks=all
# 注意 acks=all 本身不够,必须搭配陷阱 3 的 min.insync.replicas,否则会静默退化

如何避免再次触发:把 acks=all 写进生产 producer 的基线配置模板,不要让它停留在各服务自行决定。对吞吐敏感、可容忍丢失的旁路链路(如可重算的指标采样)才显式降到 acks=1,并在代码注释里写明"此链路接受丢失"——把例外变成需要主动声明的决定,而不是默认。

陷阱 2:unclean leader election 丢已提交数据

症状:一次多 broker 同时故障后集群恢复了可用性、分区重新可写,但已经确认给 producer、甚至已经被消费过的一段记录消失了。这比陷阱 1 更隐蔽——丢的是已经越过 HW、对消费者可见过的记录。

根因:unclean.leader.election.enable=true 允许一个不在 ISR 里、落后的副本当选 leader,以便在所有同步副本都不可用时恢复写入。回链 §2.3:假设 producer 已 ack 到 offset 100、HW=100,落后副本只复制到 offset 80 就被选为 leader,它会把日志截断到 80,81–100 永久丢失。这是用一致性换可用性,且丢失静默发生。

broker-server.properties Properties
# 错误写法:允许落后副本上位,截断已提交记录换取可用性
unclean.leader.election.enable=true

# 正确写法(4.x 默认即 false,但旧集群/旧 topic 可能残留 true)
unclean.leader.election.enable=false
# 代价:所有 ISR 副本都挂时分区停写,直到某个同步副本恢复——这是刻意的

如何避免再次触发:把它当成一个持久性 vs 可用性的显式声明而非性能开关。金融、订单、计费类不能丢数据的 topic 一律 false;同时在 broker 默认值与 topic 级覆盖两处都核对——一个老 topic 上残留的 true 会绕过 broker 默认。监控 UncleanLeaderElectionsPerSec 指标,任何非零都该告警。

陷阱 3:acks=all 静默退化成 acks=1(ISR 缩到 1)

症状:配置审计显示 acks=all 一切合规,但某次故障后仍然丢了已 ack 的记录。看上去和陷阱 1 矛盾——明明设了最强确认。

根因:acks=all 的语义是"所有 ISR 成员确认",不是"所有副本确认"。回链 §2.3 的关键洞察:当 follower 落后超过 replica.lag.time.max.ms 被踢出 ISR,ISR 可能收缩到只剩 leader 一台。此时 acks=all 等价于 acks=1——写到一台就算"全部 ISR 确认",下一次故障即丢。持久性的真正来源是 min.insync.replicas,不是 acks。

durability-combo.properties Properties
# 错误写法:只设 acks=all,ISR 缩到 1 时静默退化成单副本写入
acks=all
# (未设 min.insync.replicas,默认 1)

# 正确写法:broker / topic 端要求至少 2 个 ISR 成员在线
acks=all                       # producer 端
min.insync.replicas=2          # broker / topic 端
replication.factor=3           # topic 端
# 效果:ISR 缩到 1 时写入被 NotEnoughReplicas 拒绝,而不是默默欠复制

如何避免再次触发:持久性是 acks=all + min.insync.replicas=2 + replication.factor=3 三者同时成立才有的属性,缺一即降级。设计原则:min.insync.replicas = replication.factor - 1——这样允许挂一台仍可写、挂两台才停写。把"宁可拒绝写入也不欠复制"作为默认立场(fail-closed),让丢失窗口根本不存在,而不是事后补救。

反直觉提醒

acks=all 本身不够。没有 min.insync.replicas≥2,它会在 ISR 收缩时退化成 acks=1。审计配置时只看到 acks=all 就放心,是最常见的误判。

陷阱 4:replication.factor=1 与 fetch/message 尺寸错配

症状:两种形态。其一,单个 broker 下线整个分区彻底不可用、其上记录全丢。其二,producer 能成功写入大消息(broker 收下了),但这条记录始终复制不到 follower,于是它从不越过 HW、消费者永远读不到,或在故障时丢失。

根因:replication.factor=1 意味着每个分区只有一份,没有任何冗余——这通常是开发环境单节点默认值泄漏到了生产。第二种形态回链 §2.3 的复制路径:follower 像消费者一样 fetch,单次 fetch 上限由 replica.fetch.max.bytes 控制;若它小于 message.max.bytes,一条合法大消息能被 leader 接收却无法被 follower 拉取,复制就此卡死。

replication-sizing.properties Properties
# 错误写法:dev 默认漏到 prod + 复制拉取上限小于单条消息上限
replication.factor=1
replica.fetch.max.bytes=1048576    # 1 MiB
message.max.bytes=10485760         # 10 MiB —— 大消息能写入却无法被复制

# 正确写法:生产至少 3 副本,且复制拉取上限 ≥ 单条消息上限
replication.factor=3
message.max.bytes=10485760              # broker
replica.fetch.max.bytes=10485760       # >= message.max.bytes

如何避免再次触发:生产集群把 default.replication.factor=3 设为集群级默认,并禁用或受控 auto.create.topics.enable——自动建 topic 极易带出 RF=1。尺寸三件套(producer max.request.size、broker message.max.bytes、consumer fetch.max.bytes,外加 replica.fetch.max.bytes)必须当成一组联动配置统一推导,见陷阱 15。

陷阱 5:auto-commit 丢消息(提交的是 poll 返回的,不是处理过的)

症状:消费者进程崩溃或被杀后重启,发现中间一批消息从未被业务处理,但它们的 offset 已经提交、不会再投递——消息在"已消费"账面上存在,实际从未落地。

根因:enable.auto.commit=true(默认)按 auto.commit.interval.ms(默认 5s)定时提交最近一次 poll() 返回的 offset,而不是业务实际处理完的 offset。回链 §2.4 的 offset 提交语义:offset 由消费者侧维护、提交到 __consumer_offsets。设想 poll() 返回 1000 条,处理到第 501 条时定时器触发、把 offset 提交到了 1000,紧接着进程 OOM 被杀——重启从 1001 开始,502–1000 这段未处理却已提交,永久跳过。

ConsumerCommit.java Java
// 错误写法:自动提交按定时器走,提交的是 poll 返回的、不是处理完的
Properties bad = new Properties();
bad.put("enable.auto.commit", "true");      // 默认值
bad.put("auto.commit.interval.ms", "5000");
KafkaConsumer<String, String> c = new KafkaConsumer<>(bad);
while (true) {
    var records = c.poll(Duration.ofMillis(100));
    for (var r : records) process(r);       // 处理到一半崩溃 → 已提交的尾段丢失
}

// 正确写法:关掉自动提交,处理完整批后再同步提交
Properties good = new Properties();
good.put("enable.auto.commit", "false");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(good);
while (true) {
    var records = consumer.poll(Duration.ofMillis(100));
    for (var r : records) process(r);       // 先把这一批全部处理掉
    consumer.commitSync();                  // 再提交——崩溃只会重投,不会丢
}

如何避免再次触发:默认就关掉 enable.auto.commit,把 offset 提交作为业务处理成功后的显式动作。这把失败模式从"丢失"翻转成"重复"——重复可以靠消费侧幂等兜住(见陷阱 8),丢失无法补救。需要更细粒度时用 commitSync(offsets) 提交到确切处理位置。

反直觉提醒

默认 auto-commit 听起来安全,却会丢消息——它按定时器提交 poll() 返回的 offset,与业务是否处理完毫无关系。"用默认值最稳妥"在这里恰好相反。

一条记录 要安全落地 ① acks=all 堵:leader 单方面确认 ② min.insync.replicas=2 堵:ISR 缩到 1 静默退化 ③ replication.factor=3 堵:单 broker 挂掉无冗余 ④ unclean.leader.election=false 堵:落后副本上位截断 + 消费侧:处理后再提交 offset 堵:auto-commit 丢未处理段 缺一道即有缺口
图 4.1"不丢消息"是四道防线(加消费侧提交时机)叠加的组合属性,不是某个单一开关。 注意:每道防线只堵一个特定的丢失窗口(标在右侧);去掉任意一道,对应窗口就重新打开——陷阱 1–5 正是逐道防线缺失的具体后果。
想一想

RF=3、acks=all、min.insync.replicas=3。一台 broker 例行重启,producer 写入会怎样?把这个配置和"min.insync.replicas=2"对比。

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

min.insync.replicas=3 时,三个副本全部要在 ISR 里写入才成功。一台重启 → ISR 降到 2 → 不满足 3 → producer 收到 NotEnoughReplicas、分区停写,直到那台回到 ISR。即"任何一台不可用就停写"。

设成 RF-1=2 才是正解:允许挂一台仍可写、挂两台才停写,在"容忍单点故障"和"持久性"之间取平衡。把 min.insync.replicas 设成等于 RF 是常见的过度收紧——把可用性砍没了却没换来额外持久性。

4.2正确性类(重复与乱序)

这一类的共同根因是对投递语义的边界误判:把"至少一次"当成了"恰好一次",或把幂等的去重范围想得比实际大。全部回链 §2.5 投递语义。和数据丢失不同,重复与乱序能在下游补救(幂等处理、去重键),所以严重度低一档——但若误以为已经"恰好一次"而不做补救,等于把可补救的问题留成了脏数据。

陷阱 6:重试导致静默乱序

症状:单分区内消息顺序被打乱。比如同一个 key 的"创建—更新—删除"三条消息,下游收到的却是"创建—删除—更新",状态机错乱。无任何异常。

根因:max.in.flight.requests.per.connection>1(默认 5)允许同一连接上多个未确认批次并发在途。回链 §2.2 分区与顺序:顺序只在单分区内、且依赖写入到达顺序。当批次 N 失败重试、而批次 N+1 已经先落盘,重试后的 N 排到了 N+1 后面——顺序就此颠倒。

producer-ordering.properties Properties
# 错误写法:多个在途批次 + 重试 → 重试的批次可能排到后面,顺序颠倒
max.in.flight.requests.per.connection=5
retries=2147483647
# 未开幂等

# 正确写法:开启幂等,broker 按序列号保序并去重
enable.idempotence=true
# 开启后 Kafka 自动约束 in.flight<=5、acks=all、retries>0;
# 序列号机制拒绝乱序 batch,所以重试也不会打乱顺序

如何避免再次触发:enable.idempotence=true 在 4.x 已是默认,但旧配置模板或显式覆盖可能关掉它。把它当成顺序保证的开关而不只是去重开关——它通过 broker 端维护 (PID, partition) → 最高 seq、只接受 seq 恰好 +1 的 append,顺带把"重试乱序"和"重试重复"两个问题一起解决。除非有明确理由,永远不要关。

陷阱 7:ACK 丢失致重复,且幂等跨会话失效

症状:下游出现重复记录。排查发现 producer 开了幂等,单进程内确实不重复;但 producer 重启或崩溃恢复后,重复又出现了。

根因:两层。第一层,broker 写入成功但 ack 在网络上丢了,producer 超时重发——幂等 producer 靠 (PID, seq) 能识别并丢弃这个重复。第二层(边界)回链 §2.5 的 EOS 真实边界:幂等去重只在同一个 PID + 同一 producer 会话内有效。producer 重启会拿到新的 PID,broker 视之为全新生产者,旧会话发过的记录无法再去重,跨重启的重复就此漏过。

TransactionalProducer.java Java
// 错误写法:以为开了幂等就跨重启不重复——新 PID 让跨会话去重失效
props.put("enable.idempotence", "true");   // 只在单会话内去重

// 正确写法:要跨重启的恰好一次,用事务(transactional.id 靠 epoch 跨会话存活)
props.put("enable.idempotence", "true");
props.put("transactional.id", "order-writer-1");  // 稳定 ID,重启后复用
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
try {
    producer.beginTransaction();
    producer.send(record);
    // 把消费 offset 也纳入同一事务,构成消费-转换-生产闭环
    producer.sendOffsetsToTransaction(offsets, groupMetadata);
    producer.commitTransaction();
} catch (KafkaException e) {
    producer.abortTransaction();
}

如何避免再次触发:先界定需求边界。仅需"重试不重复"用幂等即可;需要跨进程重启的恰好一次,必须上事务——稳定的 transactional.id 通过 epoch 隔离僵尸实例、跨会话存活。同时记住事务的恰好一次只在 Kafka 的消费-转换-生产闭环内成立,不覆盖外部副作用(DB 写、HTTP 调用、发邮件),那些仍需自己做幂等。

反直觉提醒

幂等 producer 只在单个进程生命周期内去重。崩溃重启换了新 PID,去重立刻失效。把 enable.idempotence=true 当成端到端恰好一次,是面试和生产里最常见的过度宣称。

陷阱 8:再平衡后重复处理;"先提交后处理"反而丢失

症状:一次消费组再平衡(扩容、重启、或处理超时驱逐)之后,一批消息被重复处理了一遍——下游收到重复副作用。有人为消灭重复而改成"先提交 offset 再处理",结果重复消失了,但崩溃时变成了丢失。

根因:默认是至少一次(先处理后提交)。回链 §2.4 消费组与再平衡:消费者在处理完、提交 offset 之前崩溃或被驱逐,那批在途消息会在新 owner 上重投——这是至少一次的固有代价,正确做法是让处理幂等。反向的"先提交后处理"把语义改成了至多一次:offset 已提交但处理还没做就崩溃,那批消息再也不会投递,从重复翻转成丢失。感觉更安全的早提交,其实更危险。

IdempotentConsumer.java Java
// 错误写法:先提交后处理 → 崩溃时消息丢失(at-most-once)
for (var r : records) {
    consumer.commitSync();   // 先提交
    process(r);              // 还没处理就崩溃 → 这条永久跳过
}

// 正确写法:先处理后提交(at-least-once)+ 处理本身做成幂等
for (var r : records) {
    // 用业务唯一键去重:同一 key 重复到达时是 no-op
    upsertByBusinessKey(r.key(), r.value());   // 例如 INSERT ... ON CONFLICT DO NOTHING
}
consumer.commitSync();       // 处理完整批再提交;崩溃只会重投,幂等吸收

如何避免再次触发:接受"至少一次 + 消费侧幂等"是默认且正确的组合,不要试图靠调整提交时机消灭重复——那只会把丢失风险引进来。幂等的实现方式:业务唯一键 upsert、去重表、或下游天然幂等的操作。需要真正端到端恰好一次时回到陷阱 7 的事务方案。

反直觉提醒

"先提交后处理"把失败模式从重复翻转成丢失。重复可以靠幂等吸收,丢失不能——所以感觉更安全的早提交,才是更危险的那个。

陷阱 9:schema 演进卡死消费者

症状:producer 上线了一个新版本消息格式后,部分消费者开始大面积反序列化失败、整个消费组卡在某个 offset 推进不动,错误是 schema 不兼容。

根因:producer 发布了一个不兼容的 schema 变更(如删字段、改类型、加无默认值的必填字段)。Schema Registry 默认兼容模式是 BACKWARD——只保证新 schema 能读旧数据,且只对紧邻的上一版、非传递。在 producer 和 consumer 混版本部署、消费者读到比自己新的 schema 写的数据时,BACKWARD 兜不住。这条不直接对应 02 章某机制,根因在 schema 契约管理而非 broker——但它和顺序/投递一样属于"正确性"失败:消息在但读不出。

schema-compat.properties Properties
# 错误写法:默认 BACKWARD,混版本舰队中新写旧读会失败
# (删字段 / 改类型 / 加必填无默认值字段都会破坏兼容)

# 正确写法:混版本部署用 FULL_TRANSITIVE,并约束变更类型
# Schema Registry subject 级配置:
compatibility=FULL_TRANSITIVE
# 变更纪律:只新增可选字段 / 带默认值字段;不删、不改类型、不加必填项

如何避免再次触发:对生产者消费者混版本滚动部署的 topic,把 subject 兼容模式设为 FULL_TRANSITIVE(新旧互相可读、且对所有历史版本传递成立)。配套变更纪律:只加可选/带默认值的字段,绝不删字段、改类型或加必填项。把 schema 变更纳入 CI 校验,不兼容的变更在合并前就拦下。

4.3性能与运维类

这一类严重度最低——不丢数据、不破坏正确性,但会让吞吐归零或资源耗尽。共同根因是把并行单位(分区)和协调成本(再平衡)想得太理想,全部回链 §2.2 分区 与 §2.4 再平衡。

陷阱 10:再平衡风暴,吞吐归零

症状:消费组 lag 持续上涨、吞吐周期性归零,日志里反复出现成员离组又入组。集群没挂、消息也没丢,但消费几乎停滞。

根因:单批处理时间超过 max.poll.interval.ms(默认 5 分钟),或 session.timeout.ms 设得过紧。回链 §2.4 的关键代价:心跳跑在后台线程,但成员存活由 poll() 调用证明。处理太慢 → 两次 poll() 间隔超限 → coordinator 判定成员死亡 → 触发再平衡 → 分区重分配 → 慢消费者重新加入又再次超时 → 循环。在 eager(急切式)协议下每次再平衡都是 stop-the-world,整组停摆。

例 · 真实案例

某 240 个消费者的服务做滚动重启。重启后实例首次心跳约 38 秒到达,超过 30 秒的 session.timeout.ms,于是每个实例都先被判死、再入组,叠加滚动节奏,整组陷入约 9 分钟的再平衡 churn、消费停滞。启用 static membership 后,重启的实例凭 group.instance.id 认领回原分配、不触发再平衡——再平衡次数从 30+ 降到 0。

rebalance-tuning.properties Properties
# 错误写法:批太大处理超时 + eager 协议 + 无 static membership
max.poll.records=500
session.timeout.ms=30000           # 滚动重启首次心跳就可能超

# 正确写法:减小批量 + 协作式增量再平衡 + 静态成员
max.poll.records=100                       # 缩短单批处理时间
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
group.instance.id=consumer-instance-7      # 每实例唯一且稳定 → 重启不触发再平衡
session.timeout.ms=45000                   # 给滚动重启留出首次心跳余量

如何避免再次触发:三管齐下。① 把单批处理时间控制在 max.poll.interval.ms 内(减小 max.poll.records 或加速处理);② 用协作式增量再平衡(CooperativeStickyAssignor),只移交需要变动的分区、其余继续消费,没有全局 stop-the-world;③ 给每个实例配稳定的 group.instance.id 启用 static membership,让重启不触发再平衡。4.0 的 KIP-848 把再平衡逻辑移到 broker 端、增量收敛,进一步压缩 churn。

陷阱 11:加消费者反而 lag 增大

症状:lag 高,于是给消费组扩容;扩容后吞吐没上去,lag 反而比扩容前更大。

根因:回链 §2.4。在 eager 协议下,每加入一个新成员都触发一次 stop-the-world 再平衡——全组暂停消费、重算分配、重新拉起。如果组本身已经不稳定(处理偏慢、临近超时),新成员加入引发的再平衡会进一步挤占处理时间,把组推向陷阱 10 的风暴。另一种情形:分区数已经等于消费者数,再加的消费者分不到分区、纯闲置,却仍参与了那次再平衡的代价。

反直觉提醒

加消费者可能让 lag 更糟。每次 join 都是一次再平衡(eager 下还是 stop-the-world);分区数是消费并行度的硬上限,超过分区数的消费者只会闲置并徒增协调成本。

如何避免再次触发:扩容前先确认两件事——组是否稳定(不在临界超时),以及分区数是否还有余量(消费并行度上限 = 分区数)。先切到协作式增量再平衡再扩容,避免扩容动作本身引发 stop-the-world。若分区数已是瓶颈,问题在分区规划(陷阱 14),加消费者无效。

陷阱 12:慢消费者阻塞 poll 线程

症状:消费者偶尔因下游(数据库、HTTP)变慢而处理时间拉长,紧接着就被判死、触发再平衡,形成"下游抖动 → 再平衡 → 更慢"的连锁。

根因:在 poll() 循环里直接做重活(同步 DB 写、外部 HTTP 调用)。回链 §2.4:poll() 既拉消息又承担"成员存活"的证明,处理阻塞在 poll 线程上就等于心跳停摆。一旦单批处理超过 max.poll.interval.ms,成员被判死。

PauseResume.java Java
// 错误写法:重活直接堵在 poll 线程上,下游变慢即超时被驱逐
while (true) {
    var records = consumer.poll(Duration.ofMillis(100));
    for (var r : records) slowDbWrite(r);   // 阻塞 → 错过下一次 poll → 判死
}

// 正确写法:卸到工作线程,poll 线程用 pause/resume 维持存活而不超时
var partitions = consumer.assignment();
while (true) {
    var records = consumer.poll(Duration.ofMillis(100));
    if (!records.isEmpty()) {
        consumer.pause(partitions);          // 暂停拉取,但继续 poll 维持心跳
        submitToWorkerPool(records);         // 重活交给工作线程
    }
    if (workerPoolIdle()) {
        consumer.resume(partitions);         // 处理跟上后恢复拉取
        commitProcessedOffsets(consumer);    // 按已完成进度提交
    }
}

如何避免再次触发:保持 poll 线程轻量——只拉取和派发,重活卸到工作线程池。用 pause()/resume() 在工作线程忙时停止拉取新数据、同时继续调用 poll() 维持存活计时器,避免假死驱逐。offset 严格按工作线程实际完成的进度提交。

陷阱 13:热分区 / key 倾斜

症状:某一个消费者实例 CPU/lag 居高不下、其余实例几乎空闲,整组吞吐被这一个实例拖住。扩容无效。

根因:回链 §2.2:相同 key 经 murmur2 哈希落到同一分区以保证顺序。若 key 选择导致分布极度不均(例如按租户分区、而某个大租户占了大半流量),该租户的分区成为热点、承接它的消费者过载。Kafka 不会自动均衡倾斜的负载——分区是固定归属的。

如何避免再次触发:选高基数、分布均匀的分区 key(用 orderId 而非 userId,用 userId 而非 tenantId)。若业务要求按某低基数维度保序、又无法承受倾斜,考虑复合 key(如 tenantId + orderId)在保序粒度和均衡之间折中。规划阶段就用真实流量分布压测分区热度,别等上线发现单分区被打爆。

陷阱 14:分区数过多/过少,不可减,加分区破坏顺序

症状:三种。其一,消费并行度卡死——消费者数想加却加不上去(分区太少)。其二,controller 故障切换慢、broker 文件句柄耗尽报错(分区太多)。其三,给一个有 key 的 topic 加了分区后,下游按 key 的顺序全乱了。

根因:回链 §2.2。分区是并行的单位,消费者数不能超过分区数,过少则并行度被锁死。过多则每个 segment 占 2 个文件句柄(默认 ulimit 1024 很快耗尽)、且抬高 controller 元数据与故障切换开销(KRaft 缓解但不消除)。最隐蔽的是:分区数只能加不能减,而给有 key 的 topic 加分区会改变 key % 分区数 的映射——同一个 key 从此哈希到不同分区,历史顺序与未来顺序在不同分区之间永久错位。

反直觉提醒

分区数永远不能减,且给有 key 的 topic 加分区会静默破坏 key 顺序——murmur2 哈希的目标分区变了,同 key 消息散到新旧不同分区。分区越多也不等于越能扩展:故障切换时间、文件句柄、端到端延迟都会被推高。

如何避免再次触发:分区数当成难以回退的容量决策来定——按目标吞吐、单分区吞吐上限、未来消费并行需求保守预估,留一定余量但不盲目设大。确需调整且不能容忍顺序错乱时,新建一个分区数合适的 topic、双写或迁移,而不是在原 topic 上加分区。无 key(顺序无关)的 topic 加分区才是安全的。

陷阱 15:消息过大,三个尺寸配置不匹配

症状:producer 抛 RecordTooLargeException,或消费端拉取失败/卡住,或大消息能写入却复制不出去(与陷阱 4 的第二形态同源)。

根因:四个尺寸配置分散在三端、彼此不一致——producer 的 max.request.size、broker 的 message.max.bytes、consumer 的 fetch.max.bytes,以及复制路径的 replica.fetch.max.bytes。任意一个小于实际消息体,对应环节就断:producer 端拒发、broker 端拒收、consumer 端拉不动、follower 端复制不了。回链 §2.1 存储:记录是顺序写进 segment 的字节流,每一跳都有自己的尺寸上限。

message-size.properties Properties
# 错误写法:四个尺寸各自为政,某一跳卡住
max.request.size=1048576           # producer 1 MiB
message.max.bytes=10485760         # broker 10 MiB
fetch.max.bytes=1048576            # consumer 1 MiB —— 拉不动 broker 上的大消息

# 正确写法:四者按同一上限对齐推导(示例统一到 10 MiB)
max.request.size=10485760          # producer
message.max.bytes=10485760         # broker
replica.fetch.max.bytes=10485760   # broker 复制路径,>= message.max.bytes
fetch.max.bytes=10485760           # consumer

如何避免再次触发:把四个尺寸当成一组联动配置、从单一上限统一推导,写进配置模板而非各端独立设置。更根本的设计:大对象不要走 Kafka。Kafka 适合高吞吐的小记录;大文件用 claim-check 模式——把对象存到 S3/对象存储,消息体里只放一个指针(URL/key),消费者凭指针去取。这同时避开了尺寸配置地狱和大消息对 page cache 的污染。

想一想

一个有 key 的 topic 当前 6 分区,消费组 6 个实例刚好打满,lag 仍在涨。直接把分区加到 12 行不行?

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

不行——加分区会把 key % 6 变成 key % 12,同一个 key 从此哈希到不同分区,破坏所有下游按 key 的顺序(陷阱 14)。而且要先确认 lag 涨的根因:若是单批处理太慢或下游瓶颈,加分区/加消费者都不解决(陷阱 10、12)。

正确路径:先定位是并行度不足还是处理慢。确属分区数瓶颈、且需保序,就新建 12 分区的 topic 迁移,而不是原地加分区。若顺序无关,原地加分区才安全。

4.4反模式与"什么时候不要用 Kafka"

前面 15 条失败模式的总和,反过来勾出 Kafka 的适用边界。它是分布式可重放的分区日志——这套设计在某些场景是错的工具,硬上会持续与机制对抗。这一节服务后面的选型判断。

常见反模式

  • 把 Kafka 当请求-响应的 RPC 用:日志模型是单向追加 + 异步消费,没有内建的"响应"通道。硬做要靠关联 ID + 临时响应 topic,延迟和复杂度都远高于直接用 RPC。
  • 用消费者数超过分区数来追求并行:并行上限是分区数,多出的消费者纯闲置(陷阱 11、14)。
  • 用大量短命 topic 或海量分区做隔离:抬高 controller 负担、文件句柄、故障切换时间(陷阱 14)。
  • 把大文件直接塞进消息体:污染 page cache、撞尺寸上限——用 claim-check(陷阱 15)。
  • 靠默认配置上生产:RF=1、auto-commit、ISR 无下限都是会丢数据的默认(陷阱 3、4、5)。

什么时候不要用 Kafka

表 4.1 · Kafka 不擅长的场景与更合适的替代
场景为什么 Kafka 是错的工具更合适的方向
低吞吐、消息量小分区/副本/消费组的运维与认知成本,换不回日志模型的吞吐优势传统消息队列 / 数据库表轮询
复杂的逐条路由无内建按内容路由(topic exchange、header 路由);要在消费端自行分发RabbitMQ 等带 exchange 路由的 broker
需要内建重试 / 死信队列无原生 per-message 重试与 DLQ,需自建重试 topic + 计数RabbitMQ / SQS(原生 DLQ)
优先级队列分区日志严格按追加顺序,无消息优先级概念支持 priority 的队列系统
请求-响应 / 同步调用单向异步日志,没有响应通道gRPC / HTTP / RPC 框架
小团队、无专职运维即便 KRaft 简化了架构,集群容量规划、分区/ISR/再平衡调优仍有持续运维成本托管队列 / 云原生消息服务
高吞吐事件流 + 可重放 + 多消费组独立消费正是日志模型的主场:顺序 I/O、消费不删数据、offset 在消费侧选中 Kafka
洞察 · "更快的队列"是红旗

把 Kafka 当"更快的 MQ"会持续撞上这一章的失败模式——因为它不是队列,是日志。需要逐条路由、优先级、内建 DLQ、同步响应时,硬用 Kafka 等于跟它的设计取舍对着干。资深判断不是"什么都用 Kafka",而是认得出哪些负载与日志模型对齐、哪些不对齐。KIP-932 的 share group(4.2 GA)为部分队列语义场景提供了原生选项,但它仍架在日志之上,不改变上面这些边界。

§本章 self-check

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

  1. 一个面试官给出说法"设了 acks=all,所以不会丢消息"。指出这个说法的漏洞,并给出端到端不丢消息的完整配置组合(含消费侧)。
  2. enable.auto.commit=true 在什么具体时序下会丢消息?为什么把它关掉反而让失败模式变得"可补救"?
  3. 幂等 producer(enable.idempotence=true)能保证跨进程重启不重复吗?说明它的去重边界,以及要跨重启恰好一次该用什么。
  4. (设计题) 一个有 key 的 6 分区 topic,6 个消费者打满后 lag 仍涨。逐步说明你会怎么定位根因,以及为什么"直接加分区到 12"是错的。
答案(先做完再展开)
  1. acks=all 的语义是"所有 ISR 成员确认",ISR 缩到 1 时它退化成 acks=1(陷阱 3)。完整组合:acks=all + min.insync.replicas=2 + replication.factor=3 + unclean.leader.election.enable=false,外加消费侧处理后再提交 offset。任缺一道就有丢失窗口(图 4.1)。设计上 min.insync.replicas = RF-1 以容忍单点故障同时保持可写。
  2. auto-commit 按 auto.commit.interval.ms 定时提交最近一次 poll() 返回的 offset,与业务是否处理完无关。时序:poll 返回 1000 条,处理到第 501 条时定时器提交了 offset=1000,进程随即崩溃 → 重启从 1001 开始,502–1000 未处理却已提交、永久跳过(陷阱 5)。关掉它、改成处理完再 commitSync(),失败模式从"丢失"翻转成"重复",而重复能用消费侧幂等吸收,丢失不能。
  3. 不能。幂等去重只在同一个 PID + 同一 producer 会话内有效;进程重启拿到新 PID,旧会话的记录无法再去重(陷阱 7)。要跨重启恰好一次需用事务:稳定的 transactional.id 靠 epoch 跨会话存活、隔离僵尸实例。且事务的恰好一次只在 Kafka 消费-转换-生产闭环内成立,不覆盖 DB 写、HTTP 等外部副作用。
  4. 先定位根因而非急着扩容:① 看单个实例是否倾斜(热分区/key 倾斜,陷阱 13);② 看单批处理时间是否接近 max.poll.interval.ms 或下游是否是瓶颈(陷阱 10、12)——若是,加分区/加消费者都无效;③ 确认并行度是否真被分区数卡死。直接加分区到 12 会把 key % 6 变成 key % 12、破坏所有下游按 key 顺序(陷阱 14)。确属分区瓶颈且需保序时,新建 12 分区 topic 迁移,而非原地加分区。
进阶挑战 · 刚好够不着

给一条"绝不丢、可容忍重复"的支付事件链做端到端配置

一条支付成功事件,要求:绝不丢失、可容忍重复(下游有幂等去重)、单账户内严格有序、能容忍一台 broker 故障仍可写。写出 producer、broker/topic、consumer 三端的关键配置,并说明分区 key 怎么选、为什么不需要事务。

提示(卡住再展开)

持久性走四道防线(图 4.1);有序靠 enable.idempotence=true + 合适的分区 key(账户维度保序,但注意倾斜——见陷阱 13);"可容忍重复"意味着不需要事务,至少一次 + 消费侧幂等就够(陷阱 8)。容忍一台故障 ⇒ min.insync.replicas = RF-1 = 2。逐项对照陷阱 1/2/3/6/8 的"正确写法"。