这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 “Kafka 会不会丢消息?”的正确回答是:配置得当就不丢,配置不当每一层都能丢。本章按端到端路径给出工程清单。
18.1 可靠性是端到端属性
业务库 -> Producer -> 网络 -> Broker集群 -> 消费者 -> 下游存储
任何一环出问题都可能丢:
- 业务代码发送前进程崩溃;
- Producer 异步缓冲未 flush 就退出;
acks=0/1且 Leader 副本未同步; - ISR 收缩到 1 时发生 unclean 选举;
- 消费者先提交位移后处理;
- 消费者 auto.offset.reset=latest 且位移过期。
下面逐层堵漏。
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:
- 一台宕机:ISR 剩 2,仍满足 min ISR,写入继续且已提交消息双副本持久;
- 两台同时宕机:写入拒绝(
NotEnoughReplicasException),但不丢; - 若开 unclean election:落后副本上位会丢已提交数据,与“不丢”目标冲突。
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 -> 跳过所有堆积消息。
防护:
- 关键 topic 的保留期大于最长可接受停机时间;
- 严格监控 lag,不允许长期堆积;
auto.offset.reset明确设置并写入团队规范;- 长期不消费的组显式记录,避免“复活时踩坑”。
18.6 顺序与重复的权衡
- 要不丢:至少一次 + 幂等;
- 要顺序:同 key 同分区 + 幂等生产者 + in-flight<=5;
- 要精确一次(Kafka 内):事务 + read_committed;
- 涉及外部系统:本地事务 + 去重表。
不要试图用“关闭重试”换取不重复——那是用丢消息换不重复,方向反了。
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 告警
复盘要点:不只修配置,还要让“同类错误无法再次发生”(配置准入、告警、演练)。
本章小结
- 不丢消息 = 生产端可靠发送 + Broker 多副本 + 消费端正确提交,缺一不可;
- 黄金组合:
acks=all + RF=3 + min.insync.replicas=2 + unclean=false; - 消费端铁律:先处理后提交,失败必须显式处理(重试/死信);
- 位移过期 + latest 是隐蔽的丢失来源;
- 把清单工具化:CI 校验配置、告警兜底、事故复盘闭环。
思考题
- 每一层都“看似可靠”,为什么拼起来仍可能丢消息?举一个跨层时序的例子。
- 为什么
min.insync.replicas=2在两台 Broker 宕机时反而“保护”了你? - 为你的系统写一份“消息可靠性 SLA”:允许丢吗?允许重吗?允许乱序吗?