一句话总结
消息可靠性有三个层次:不丢(At-least-once)、不重(Exactly-once)、有序.它们不是 Kafka 单方面保证的,而是 Producer 配置 + Broker 配置 + Consumer 处理逻辑三方配合的结果.
本篇把端到端链路拆开,逐段分析哪里会丢、哪里会重、怎么保证顺序.
消息传递语义:三种级别
| 语义 | 含义 | 代价 |
|---|---|---|
| At-most-once | 最多投递一次,可能丢,不会重 | 最简单,性能最高 |
| At-least-once | 至少投递一次,不会丢,可能重 | 需要重试 + Consumer 幂等 |
| Exactly-once | 恰好一次,不丢不重 | 最复杂,需要事务或幂等机制 |
Kafka 默认配置(franz-go)实现的是 At-least-once:Producer 幂等保证不丢,但 Consumer 侧可能重复处理.要做到端到端 Exactly-once,需要额外机制.
第一段:Producer → Broker(不丢)
会丢的场景
| 场景 | 原因 | 解法 |
|---|---|---|
| acks=0 | 发出去不等回复,网络丢包就没了 | 设 acks=all |
| acks=1 + Leader 宕机 | Leader 写入后还没同步到 Follower 就挂了 | 设 acks=all + min.insync.replicas≥2 |
| Producer 内存缓冲区断电 | 消息还没发出去,进程挂了 | Outbox 模式(上游先落库) |
| 网络超时 + 不重试 | 消息发了但 ack 丢失,Producer 认为失败却不重试 | 开启重试(franz-go 默认无限重试可恢复错误) |
不丢的标准配置
|
Broker 端:
|
三者配合的效果:
|
第二段:Broker 存储(不丢)
消息写入 Broker 后是否安全?取决于:
| 风险 | 原因 | Kafka 的保护 |
|---|---|---|
| 单机磁盘损坏 | 硬件故障 | 副本机制:数据在多个 Broker 上有拷贝 |
| 所有副本同时丢失 | 机房级故障 | 跨机架部署(broker.rack配置) + 异地多集群(MirrorMaker) |
| 消息过期被删 | retention 到期 | 调大 retention 或用 compact 策略 |
只要副本因子 ≥ 2 且副本分布在不同故障域,Broker 端的数据丢失概率极低.
第三段:Broker → Consumer(不丢 + 可能重复)
这是最容易出问题的一段.Consumer 拉到消息后,有一个处理窗口:
|
不丢:先处理后提交
|
重复的根源
"先处理后提交"保证了不丢,但带来了重复:
|
这就是 At-least-once 的本质:不丢和不重在"先处理后提交"模型下不可兼得,除非引入额外机制.
消费端去重:幂等处理
既然 At-least-once 下消息可能重复到达,Consumer 侧需要幂等(idempotent)处理——同一条消息处理多次和处理一次的最终效果一样.
策略一:业务天然幂等
有些操作天然幂等,不需要额外处理:
|
策略二:唯一键去重
|
策略三:业务操作与 offset 在同一事务
|
这是真正的 Consumer 端 Exactly-once:即使重复消费,数据库里的 offset 已经推进过了,事务会被忽略.
端到端 Exactly-once:Kafka 事务
Kafka 0.11+ 提供了事务 API,支持跨分区原子写入 + Consumer offset 原子提交:
|
这是 Kafka Streams(Java)的核心能力.在 Go 中,franz-go 支持事务 API:
|
事务的适用场景
Kafka 事务解决的是 consume-transform-produce 模式的 EOS(Exactly-Once Semantics).
如果你的 Consumer 是写外部系统(数据库、Redis、ES),事务帮不了你——因为外部系统不参与 Kafka 事务.此时回退到"幂等处理"策略.
顺序性保证
顺序性在前面篇幅已多次提到,这里做一个完整总结:
Producer 端顺序
| 条件 | 顺序保证 |
|---|---|
| 同一 Key,幂等 Producer | ✓ 分区内严格有序(Broker 按 sequence number 排序) |
同一 Key,非幂等 + max.in.flight > 1 |
✗ 重试可能乱序 |
| 不同 Key | ✗ 落入不同 Partition,无跨分区顺序 |
Consumer 端顺序
| 条件 | 顺序保证 |
|---|---|
| 单 Consumer 串行处理 | ✓ 按 offset 顺序处理 |
| 按 Partition 分 goroutine,Partition 内串行 | ✓ 每个 Partition 内有序 |
| Worker Pool 打散处理 | ✗ 无序(需要业务容忍或自行排序) |
端到端有序的完整条件
|
满足以上三点,同一个 Key 的消息从生产到消费严格有序.
总结:可靠性配置速查
| 目标 | Producer 配置 | Broker 配置 | Consumer 策略 |
|---|---|---|---|
| 不丢 | acks=all, 幂等开启 | min.insync.replicas=2, replication.factor=3 | 先处理后提交 |
| 不重 | (Producer 幂等保证 Broker 端不重) | — | 业务幂等 / 唯一键去重 / 事务 |
| 有序 | 设 Key + 幂等 | — | Partition 内串行处理 |
| 端到端 EOS | TransactionalID | transaction.state.log 配置 | 事务消费(仅 consume-produce 模式) |
下一篇搭建本地开发环境:docker-compose 起 3-Broker 集群 + kafka-ui,常用 CLI 命令,监控指标入门.