KafkaNotes

第 18 章:可靠性最佳实践:不丢消息的完整清单

zjc 于 2026-01-18 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 “Kafka 会不会丢消息?”的正确回答是:配置得当就不丢,配置不当每一层都能丢。本章按端到端路径给出工程清单。

18.1 可靠性是端到端属性

业务库 -> Producer -> 网络 -> Broker集群 -> 消费者 -> 下游存储

任何一环出问题都可能丢:

下面逐层堵漏。

18.2 生产端清单

检查项 正确做法
发送确认 acks=all
重试 enable.idempotence=true(含自动重试)
超时预算 delivery.timeout.ms 覆盖重试总时长,按业务 SLA 设置
回调 必须实现,失败落补偿表/告警,禁止静默吞掉
退出 close()flush(),确保缓冲区清空
消息大小 不超过 max.request.size 与 Broker message.max.bytes
本地事务 业务库 + 消息发送用 Outbox 模式(第 26 章)

典型错误代码:

producer.send(record); // 无回调无检查
System.exit(0);       // 缓冲区消息可能尚未发出

正确姿势:

producer.send(record, (md, ex) -> {
    if (ex != null) saveToOutboxRetry(record, ex); // 或直接本地补偿
});
producer.flush();

18.3 Broker 端清单

检查项 推荐配置
副本因子 replication.factor=3(重要业务)
最小同步副本 min.insync.replicas=2
脏选举 unclean.leader.election.enable=false
机架分布 broker.rack 设置,副本跨机架/可用区
内部主题 __consumer_offsets RF=3
磁盘 RAID10 或多副本替代;监控坏盘
保留期 覆盖最长故障恢复时间,防止位移过期

为什么是 3 + 2 + false

18.4 消费端清单

检查项 正确做法
位移提交 手动提交,先处理后提交
提交失败 commitSync 捕获后重试/记录,必要时暂停消费
再平衡 onPartitionsRevoked 中同步提交
起点策略 明确 auto.offset.reset(重要业务用 earliest)
处理失败 不 ack,交给重试/死信,禁止 catch 后吞掉
幂等 唯一键/去重表/状态机,容忍重复
堆积 监控 lag,防止位移超过保留期被删除

危险代码:

for (var r : records) {
    try { process(r); }
    catch (Exception e) { log.error("ignore", e); } // 吞掉 = 消息丢失
}
consumer.commitSync();

失败消息必须进入显式的重试/死信流程,并保留可观测性。

18.5 位移过期:容易被忽视的丢失

场景:消费者组故障超过 7 天(默认保留期),位移记录被删除,重启后 auto.offset.reset=latest -> 跳过所有堆积消息

防护:

18.6 顺序与重复的权衡

不要试图用“关闭重试”换取不重复——那是用丢消息换不重复,方向反了。

18.7 上线前检查清单

Topic:
  [ ] replication.factor = 3
  [ ] min.insync.replicas = 2
  [ ] retention 覆盖故障恢复窗口
  [ ] 分区数满足吞吐与并行需求

Producer:
  [ ] acks=all, idempotence=true
  [ ] 发送失败有回调与补偿
  [ ] 停机 flush/close
  [ ] 压缩与大消息策略确认

Consumer:
  [ ] 手动提交,先处理后提交
  [ ] auto.offset.reset 明确
  [ ] 失败进重试/死信
  [ ] 幂等设计评审
  [ ] rebalance 参数与静态成员确认

运维:
  [ ] unclean.leader.election.enable=false
  [ ] 副本跨机架
  [ ] lag/ISR/Controller 告警
  [ ] 容量与磁盘增长预警

18.8 一个真实事故复盘模板

现象:   支付回调消息丢失 3000 条
时间线: 14:00 Broker2 宕机 -> 14:05 恢复 -> 15:00 业务发现
根因:   topic RF=1,Broker2 上的分区数据未同步副本
修复:   RF 扩到 3 + min ISR 2,补发丢失区间
改进:   上线 topic 配置校验(禁止 RF<3),增加 ISR 告警

复盘要点:不只修配置,还要让“同类错误无法再次发生”(配置准入、告警、演练)。

本章小结

思考题

  1. 每一层都“看似可靠”,为什么拼起来仍可能丢消息?举一个跨层时序的例子。
  2. 为什么 min.insync.replicas=2 在两台 Broker 宕机时反而“保护”了你?
  3. 为你的系统写一份“消息可靠性 SLA”:允许丢吗?允许重吗?允许乱序吗?