跳过导航

Kafka 消息为什么仍可能重复或丢失?生产、Broker、消费三段治理

约 7 分钟...次浏览
专栏分布式系统与消息队列第 2 篇

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.idempotenceacksretries 已有更安全的默认组合;这里显式写出是为了把可靠性合同固化在配置与测试中。通常不要靠手工调小 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.recordsmax.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 保存数据,消费端选择业务成功后提交并接受重复,再用唯一键或状态机实现幂等。可靠系统的目标不是假装故障不存在,而是让每一次不确定都可以重试、识别和恢复。

分享:
文章作者:狼码纪
版权声明:本博客所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。文章可能参考了其他优秀文章,如有侵权请联系删除。