Kafka 消息为什么仍可能重复或丢失?生产、Broker、消费三段治理
Kafka 被认为是可靠消息系统,但“使用 Kafka”并不自动等于“不丢不重”。只要系统存在网络超时、进程崩溃和重试,就会出现一个经典难题:发送方不知道对方究竟没处理,还是已经处理但确认丢了。
可靠性必须分成三段分析:生产者到 Broker、Broker 内部副本、消费者到业务结果。
一、生产端为什么会丢
若生产者使用 acks=0,消息发出后不等待 Broker 确认,网络中断或 Broker 拒绝时客户端也不知道。acks=1 只要求 Leader 写入成功,Leader 在副本同步前故障仍存在丢失窗口。
推荐的可靠性基线:
acks=all
enable.idempotence=true
retries=2147483647
delivery.timeout.ms=120000
在 Kafka 3.x/4.x 客户端中,enable.idempotence、acks 和 retries 已有更安全的默认组合;这里显式写出是为了把可靠性合同固化在配置与测试中。通常不要靠手工调小 retries 控制总时长,而应以 delivery.timeout.ms 约束一次 send() 从入队到最终成功/失败的总预算,并保证它不小于 request.timeout.ms + linger.ms。
发送必须检查异步回调:
kafkaTemplate.send("orders", orderId, event)
.whenComplete((result, ex) -> {
if (ex != null) {
failedMessageStore.save(event, ex.getMessage());
}
});
只调用 send() 却忽略 Future,并不代表消息发送成功。
二、重试为什么产生重复
生产者发出消息后,Broker 可能已经写入,但确认响应在网络中丢失。生产者超时重试,Broker 若无法识别这是同一次发送,就会追加第二条相同消息。
幂等生产者通过 Producer ID 与分区序列号识别单个生产者会话内的重复批次。它能解决重试导致的 Kafka 日志重复,但不能替代业务幂等:应用重启后重新构造同一个订单事件、数据库轮询重复发布,仍可能产生语义重复。
三、Broker 副本层的可靠性
Topic 应配置合理副本,并约束 ISR:
replication.factor=3
min.insync.replicas=2
producer acks=all
这表示 ISR 数量少于 2 时拒绝写入;正常 acks=all 会等待当时 ISR 中的全部副本确认,而非固定只等两个。副本必须跨 Broker、机架或可用区合理放置,否则“三副本”也可能共享同一故障域。
若允许传统意义上的不干净选主,落后且不安全的副本可能成为 Leader,已经确认但未复制给它的数据会丢失。生产环境通常应禁止这种数据丢失式选主,除非业务明确选择“可用性优先、允许数据缺口”。Kafka 4.x KRaft 的 Eligible Leader Replicas 可以记录虽暂时不在 ISR、但被控制器判定为安全的候选者,不应把这种安全选主与 unclean.leader.election.enable=true 混为一谈。
四、消费端最常见的丢失窗口
先提交 offset,再执行业务:
poll -> commit offset -> 写数据库 -> 进程崩溃
若进程在提交后、业务完成前崩溃,重启后消费者从新 offset 开始,这条消息不会再处理,形成业务丢失。
自动提交是在 poll() 驱动下按周期提交先前返回批次的位置;若业务被异步丢进线程池、批次尚未真正完成就再次 poll,它就会与实际处理进度脱节。更安全的顺序通常是业务成功后再提交 offset:
poll -> 处理业务成功 -> commit offset
五、业务成功后提交为什么会重复
写数据库成功 -> 进程崩溃 -> offset 尚未提交
重启后消息再次投递。数据库已完成第一次处理,于是重复扣款、重复发券等问题出现。
因此 at-least-once 消费通常选择“不轻易丢,但允许重复”,再由业务幂等消化重复。
六、业务幂等的几种实现
唯一业务键
CREATE UNIQUE INDEX uk_payment_event
ON payment_record(event_id);
消费事务中先插入事件处理记录,唯一键冲突代表已经处理。业务更新与去重记录必须在同一个本地事务中,否则仍有状态不一致窗口。
状态机条件更新
UPDATE orders
SET status = 'PAID'
WHERE id = ? AND status = 'UNPAID';
通过受影响行数判断是否首次完成状态迁移。
Inbox 表
保存 event_id、消费者名称、处理结果与时间,适合审计和重放。需要清理策略,否则表会无限增长。
Redis SETNX 可做快速挡板,但若 Redis 去重成功、数据库业务失败,就可能永久跳过消息。它不能简单替代同库事务内的唯一约束。
七、序列化失败与“毒消息”
消息成功进入 Kafka,不代表消费者一定能反序列化。Schema 不兼容、类名变更或脏数据可能让消费线程反复失败,后续消息被阻塞。
应设置错误处理、有限重试和死信主题,并保留原 Topic、Partition、Offset、异常类型和原始 payload,方便修复后重放。
八、再均衡造成的重复窗口
消费者处理批次过慢、心跳异常或扩缩容会触发 rebalance。分区被回收时,如果实际处理进度尚未正确提交,新消费者会从旧 offset 重新读取。
要控制单批大小与处理时长,合理设置 max.poll.records 和 max.poll.interval.ms,并在分区撤销回调中谨慎提交已完成进度。
九、端到端可靠性设计
一个常见方案是:
数据库本地事务写业务数据 + Outbox
-> 发布器至少一次投递 Kafka
-> Kafka 三副本、acks=all
-> 消费者至少一次消费
-> 数据库唯一键/状态机实现幂等
-> 成功后提交 offset
这套方案允许重复,但每个故障窗口都可恢复,比依赖一次“完美网络调用”可靠得多。
十、排查清单
- Producer 是否等待
acks=all并检查发送结果? - 是否启用幂等生产,重试和超时是否合理?
- Topic 副本数、ISR 与不干净选主配置是否符合目标?
- offset 是在业务前还是业务后提交?
- 消费业务是否有数据库级唯一键或条件更新?
- 去重记录和业务变更是否处于同一事务?
- 是否存在反序列化失败、无限重试或 rebalance?
- 是否在升级到 Kafka 4.x 时验证过 KRaft、生产者冲突配置和新再均衡协议的行为?
- 能否根据 eventId、topic、partition、offset 追踪一条消息?
参考基线
总结
消息丢失与重复不是 Kafka 单点参数问题,而是端到端确认边界问题。生产端用确认和幂等降低写入风险,Broker 用副本与 ISR 保存数据,消费端选择业务成功后提交并接受重复,再用唯一键或状态机实现幂等。可靠系统的目标不是假装故障不存在,而是让每一次不确定都可以重试、识别和恢复。
相关文章
Exactly Once 是真实能力还是营销概念?先定义处理边界
区分投递一次、处理一次与业务效果一次,分析 Kafka 幂等生产者、事务、read-process-write 链路及外部数据库边界,说明 Exactly Once 的真实适用范围。
消息顺序性到底应该如何保证?从 Kafka 分区到业务状态机
区分全局有序、分区有序与业务实体有序,分析生产并发、分区映射、消费者并行、失败重试和扩分区导致的乱序,并给出工程方案。
Kafka 为什么能实现高吞吐?顺序写、批处理、页缓存与零拷贝
从分区日志、顺序追加、批量压缩、操作系统页缓存、sendfile、消费者拉取和横向扩展解释 Kafka 高吞吐的完整机制。