消息队列中的幂等:从重复消费到 Kafka Exactly Once
消息队列中的幂等:从重复消费到 Kafka Exactly Once
在消息队列中,一个核心的问题是:
同一条消息可能会被重复投递、重复消费,系统应该如何避免业务结果被重复执行?
什么是幂等
所谓幂等,可以简单理解为:
同一个操作执行一次和执行多次,最终结果相同。
例如:
1 | |
执行一次和执行十次,最终都是:
1 | |
这是天然幂等的。
但下面这种操作不是幂等的:
1 | |
如果执行两次,最终结果就和执行一次不同。
因此,在消息队列中,只要 Consumer 执行的是这类“累加、扣减、创建”操作,就必须考虑重复消费的问题。
为什么 Consumer 会重复消费
很多消息队列采用 At-Least-Once,也就是“至少一次”投递语义。
例如:
1 | |
Broker 没收到 ACK,于是认为这条消息可能没有成功处理。
因此重新投递:
1 | |
但实际上第一次扣库存已经成功。
如果 Consumer 再执行:
1 | |
那么库存就被扣了两次。
所以:
At-Least-Once 保证消息尽量不丢,但代价是可能重复消费。
这也是 Consumer 幂等存在的原因。
message_id:给逻辑消息一个唯一身份
一种常见方案,是给每条逻辑消息一个唯一的 message_id。
例如:
1 | |
注意:
一条逻辑消息可以被投递很多次,但
message_id不应该改变。
例如:
1 | |
而不能第二次重新生成一个新的 ID。
因此 message_id 标识的是“逻辑消息”而不是“一次网络发送”
Consumer 如何利用 message_id 实现幂等
Consumer 可以维护一张“已经处理过的消息”表:
1 | |
消费时:
1 | |
这样即使 MQ 重复投递:
1 | |
业务也只执行一次。
message_id 通常是谁生成的
通常是 Producer 在创建逻辑消息时生成。
例如:
1 | |
可以使用:
1 | |
例如:
1 | |
但业务 ID 不能随便直接拿来使用。
例如同一个订单可能产生:
1 | |
所以只使用:
order_id = 1001
显然不能区分这些不同事件。
因此常见做法是:
1 | |
或者直接给每个事件一个独立 message_id。
数据库唯一约束也可以做幂等
实际业务中,不一定非要维护独立的 processed_messages 表。
例如“一个订单支付后只能发一次奖励”。
数据库可以设置:
1 | |
那么:
1 | |
只能存在一条。
MQ 即使重复投递:
1 | |
第二次插入时数据库也会拒绝重复数据。
这比简单写:
1 | |
更可靠。
因为两个 Consumer 可能并发执行:
1 | |
然后同时插入。
而数据库的 UNIQUE 约束可以从存储层保证最终只有一个成功。
幂等记录和业务操作最好放在一个事务中
下面这种写法仍然有问题:
1 | |
如果在第 3 步崩溃:
1 | |
MQ 重投以后,积分又会增加一次。
因此应该:
1 | |
让它们一起成功或一起失败
也就是说:
Consumer 幂等通常需要“唯一标识 + 原子事务”配合。
Producer 也存在幂等问题
重复问题不仅发生在 Consumer 端。
例如:
1 | |
如果 Broker 再写一次:
1 | |
那么消息本身就在 MQ 中重复了。
因此还需要 Producer 侧的幂等。
Kafka 的幂等 Producer
Kafka 并不是简单给每条消息生成一个全局 UUID。
Kafka 会给 Producer 一个身份:
1 | |
然后对每个 Partition 中发送的消息维护递增的:
1 | |
例如:
1 | |
如果:
1 | |
第一次已经被 Broker 写入,但 ACK 丢失,Producer 重试:
1 | |
Broker 可以识别:
这条消息已经写过了。
因此不会再次追加。
所以 Kafka 的幂等 Producer 本质上是:
1 | |
实现重复发送检测。
Kafka 的 Exactly Once
仅仅 Producer 幂等还不够。
假设:
1 | |
Consumer:
1 | |
如果第 3 步完成后,在第 4 步之前崩溃:
1 | |
Consumer 重启后会再次读取同一条消息,于是再次写 B。
Kafka 使用事务解决这个问题:
1 | |
因此:输出结果 + 消费进度 成为一个原子操作。
要么全部成功,要么全部失败。
所以可以把 Kafka EOS 简化为:
1 | |
Kafka 的 Exactly Once 不是万能的
Kafka 的 Exactly Once 最适合:
1 | |
但如果是:
1 | |
Kafka 无法自动保证 MySQL 中:
1 | |
只执行一次。
例如:
1 | |
Consumer 重启后仍然会再次处理这条消息。
此时依然需要:
1 | |
因此:
Exactly Once 通常不是 MQ 单独就能完全保证的,而是一个端到端系统语义。
总结
消息队列中的幂等问题,本质上来自:
网络不可靠
ACK 可能丢失
进程可能崩溃
消息可能重试
因此系统需要接受一个事实:
“重复”很难完全避免,但可以让重复不产生额外副作用。
Consumer 侧通常通过:
message_id + 数据库唯一约束 + 事务
保证业务幂等
Kafka Producer 侧则通过:
1 | |
避免重试导致消息重复写入。
Kafka Exactly Once 再进一步通过:
1 | |
保证 Kafka 内部流处理链路中的 Exactly Once。
Consumer 幂等解决的是“同一条消息被重复处理怎么办”。
Producer 幂等解决的是“同一次发送因为重试被 Broker 写入多次怎么办”。
而 Exactly Once,则是在这两类问题之上进一步解决:
消费进度和处理结果如何保证原子一致。