先说结论

消息系统常见语义包括最多一次、至少一次和恰好一次。工程上最常见的是“至少一次投递 + 业务幂等”,因为跨 Kafka、数据库和外部服务的端到端恰好一次需要所有环节共同参与。

生产端

acks=all 要求 Leader 等待同步副本确认,可靠性高于只等待 Leader。它需要与合理的副本数、min.insync.replicas 配合,否则副本不足时仍可能降低容灾能力或导致写入失败。

生产者应开启幂等能力,使用有限重试并监控最终失败。幂等生产者可以避免重试在单分区内产生重复写入,但不能自动让下游数据库操作变成幂等。

Broker 端

建议生产 Topic 使用多个副本,并将副本分布在不同故障域。同步副本集合 ISR 反映当前能跟上 Leader 的副本。若允许不同步副本被选为 Leader,可能以数据丢失换取可用性,需要明确取舍。

仅配置副本数还不够,还要监控 ISR 收缩、离线分区、磁盘容量和请求延迟,并定期演练 Broker 故障。

消费端位点提交

如果先提交位点再处理消息,处理失败会丢失业务结果;如果处理成功后再提交,提交失败会导致消息再次消费。因此通常选择后者,并让业务处理支持幂等。

拉取消息 → 执行业务事务 → 成功后提交位点

关闭自动提交可以更准确地控制时机,但需要正确处理批量消息、异常、再均衡和进程退出。不要在一批消息只处理一半时直接提交整批最大位点。

业务幂等方案

  • 使用业务唯一键建立数据库唯一约束。
  • 保存 topic + partition + offset 作为消费记录。
  • 状态更新使用条件版本号,防止重复推进。
  • 调用外部接口时传递幂等键。

“先查是否处理过,再执行插入”在并发下仍可能竞争,数据库唯一约束或原子操作才是最终防线。

数据库与 Kafka 的一致性

业务写数据库后再发 Kafka,任一步骤失败都会不一致。常见解决方案是 Outbox:在同一个数据库事务中写业务表和事件表,再由独立任务可靠投递事件。消费者侧仍需要幂等。

Kafka 事务能原子地写入多个 Kafka 分区以及提交消费位点,适合 Kafka-to-Kafka 流程,但无法直接替代与普通关系数据库之间的分布式一致性设计。

排查清单

  1. 生产者最终发送失败是否被记录和告警。
  2. Topic 副本数、ISR 和最小同步副本是否合理。
  3. 消费位点是在业务成功之前还是之后提交。
  4. 消费逻辑是否拥有真正的原子幂等约束。
  5. 再均衡、超时和进程崩溃场景是否经过故障演练。

常见问题

追问:acks=all 就绝对不丢吗?

不是。还取决于副本数、ISR、最小同步副本、磁盘故障、生产者错误处理和运维配置。它只是端到端可靠链路中的一环。

追问:Kafka 的 exactly-once 等于业务恰好一次吗?

不等于。Kafka 事务适合 Kafka 内部的读写链路;写普通数据库或调用外部接口仍需幂等与一致性方案。