消息队列中的幂等:从重复消费到 Kafka Exactly Once

消息队列中的幂等:从重复消费到 Kafka Exactly Once

在消息队列中,一个核心的问题是:

同一条消息可能会被重复投递、重复消费,系统应该如何避免业务结果被重复执行?

什么是幂等

所谓幂等,可以简单理解为:

同一个操作执行一次和执行多次,最终结果相同。

例如:

1
status = PAID

执行一次和执行十次,最终都是:

1
PAID

这是天然幂等的。

但下面这种操作不是幂等的:

1
2
3
points += 100
stock -= 1
balance += 100

如果执行两次,最终结果就和执行一次不同。

因此,在消息队列中,只要 Consumer 执行的是这类“累加、扣减、创建”操作,就必须考虑重复消费的问题。


为什么 Consumer 会重复消费

很多消息队列采用 At-Least-Once,也就是“至少一次”投递语义。

例如:

1
2
3
4
5
6
7
Broker -> Consumer: msg100

Consumer:
1. 扣库存成功
2. 准备 ACK

此时 Consumer 崩溃

Broker 没收到 ACK,于是认为这条消息可能没有成功处理。

因此重新投递:

1
Broker -> Consumer: msg100

但实际上第一次扣库存已经成功。

如果 Consumer 再执行:

1
stock -= 1

那么库存就被扣了两次。

所以:

At-Least-Once 保证消息尽量不丢,但代价是可能重复消费。

这也是 Consumer 幂等存在的原因。


message_id:给逻辑消息一个唯一身份

一种常见方案,是给每条逻辑消息一个唯一的 message_id

例如:

1
2
3
4
5
{
"message_id": "abc-123",
"type": "OrderPaid",
"order_id": 1001
}

注意:

一条逻辑消息可以被投递很多次,但 message_id 不应该改变。

例如:

1
2
3
4
5
6
7
第一次投递:
message_id = abc-123

ACK 丢失

第二次重投:
message_id = abc-123

而不能第二次重新生成一个新的 ID。

因此 message_id 标识的是“逻辑消息”而不是“一次网络发送”

Consumer 如何利用 message_id 实现幂等

Consumer 可以维护一张“已经处理过的消息”表:

1
2
3
4
5
6
processed_messages

message_id
----------------
abc-123
def-456

消费时:

1
2
3
4
5
6
7
8
9
收到 abc-123
|
v
查询是否处理过
|
+-- 已处理 -> 直接跳过
|
+-- 未处理 -> 执行业务
记录 abc-123

这样即使 MQ 重复投递:

1
2
3
abc-123
abc-123
abc-123

业务也只执行一次。


message_id 通常是谁生成的

通常是 Producer 在创建逻辑消息时生成。

例如:

1
2
3
4
5
6
7
业务事件产生
|
v
生成 message_id
|
v
发送消息

可以使用:

1
2
3
UUID
Snowflake ID
业务唯一键

例如:

1
OrderPaid:1001

但业务 ID 不能随便直接拿来使用。

例如同一个订单可能产生:

1
2
3
OrderCreated:1001
OrderPaid:1001
OrderShipped:1001

所以只使用:
order_id = 1001
显然不能区分这些不同事件。

因此常见做法是:

1
event_type + business_id

或者直接给每个事件一个独立 message_id


数据库唯一约束也可以做幂等

实际业务中,不一定非要维护独立的 processed_messages 表。

例如“一个订单支付后只能发一次奖励”。

数据库可以设置:

1
UNIQUE(order_id, reward_type)

那么:

1
1001 | ORDER_PAID

只能存在一条。

MQ 即使重复投递:

1
2
OrderPaid(1001)
OrderPaid(1001)

第二次插入时数据库也会拒绝重复数据。

这比简单写:

1
2
3
if (!exists()) {
insert();
}

更可靠。

因为两个 Consumer 可能并发执行:

1
2
Consumer A: 不存在
Consumer B: 不存在

然后同时插入。

而数据库的 UNIQUE 约束可以从存储层保证最终只有一个成功。


幂等记录和业务操作最好放在一个事务中

下面这种写法仍然有问题:

1
2
3
4
1. 检查 message_id
2. 加积分
3. 程序崩溃
4. 记录 message_id

如果在第 3 步崩溃:

1
2
积分已经加了
但 message_id 没记录

MQ 重投以后,积分又会增加一次。

因此应该:

1
2
3
4
5
6
7
BEGIN

记录 message_id
+
执行业务操作

COMMIT

让它们一起成功或一起失败

也就是说:

Consumer 幂等通常需要“唯一标识 + 原子事务”配合。


Producer 也存在幂等问题

重复问题不仅发生在 Consumer 端。

例如:

1
2
3
4
5
6
7
Producer -> Broker

Broker 已经成功写入
ACK 丢失

Producer 认为发送失败
于是重试

如果 Broker 再写一次:

1
2
msg
msg

那么消息本身就在 MQ 中重复了。

因此还需要 Producer 侧的幂等。


Kafka 的幂等 Producer

Kafka 并不是简单给每条消息生成一个全局 UUID。

Kafka 会给 Producer 一个身份:

1
PID

然后对每个 Partition 中发送的消息维护递增的:

1
Sequence Number

例如:

1
2
3
4
5
6
PID = 42

Partition 0:
seq = 0
seq = 1
seq = 2

如果:

1
PID=42, seq=2

第一次已经被 Broker 写入,但 ACK 丢失,Producer 重试:

1
PID=42, seq=2

Broker 可以识别:

这条消息已经写过了。

因此不会再次追加。

所以 Kafka 的幂等 Producer 本质上是:

1
2
3
4
5
PID
+
Partition
+
Sequence Number

实现重复发送检测。


Kafka 的 Exactly Once

仅仅 Producer 幂等还不够。

假设:

1
2
3
4
5
6
Topic A
|
Consumer
|
v
Topic B

Consumer:

1
2
3
4
1. 读取 A
2. 处理
3. 写入 B
4. 提交 A 的 offset

如果第 3 步完成后,在第 4 步之前崩溃:

1
2
Topic B 已经写入
但 offset 没提交

Consumer 重启后会再次读取同一条消息,于是再次写 B。

Kafka 使用事务解决这个问题:

1
2
3
4
5
6
7
BEGIN TRANSACTION

写 Topic B
+
提交 Topic A 的 Consumer Offset

COMMIT

因此:输出结果 + 消费进度 成为一个原子操作。

要么全部成功,要么全部失败。

所以可以把 Kafka EOS 简化为:

1
2
3
4
5
幂等 Producer
+
Kafka Transaction
+
Offset 与输出结果一起提交

Kafka 的 Exactly Once 不是万能的

Kafka 的 Exactly Once 最适合:

1
Kafka -> 处理 -> Kafka

但如果是:

1
2
3
4
5
Kafka
|
Consumer
|
MySQL

Kafka 无法自动保证 MySQL 中:

1
balance += 100

只执行一次。

例如:

1
2
3
MySQL 修改成功
Consumer 崩溃
Kafka offset 没提交

Consumer 重启后仍然会再次处理这条消息。

此时依然需要:

1
2
3
4
5
业务幂等
UNIQUE 约束
message_id
数据库事务
Outbox 等机制

因此:

Exactly Once 通常不是 MQ 单独就能完全保证的,而是一个端到端系统语义。


总结

消息队列中的幂等问题,本质上来自:

  • 网络不可靠

  • ACK 可能丢失

  • 进程可能崩溃

  • 消息可能重试

因此系统需要接受一个事实:

“重复”很难完全避免,但可以让重复不产生额外副作用。

Consumer 侧通常通过:
message_id + 数据库唯一约束 + 事务
保证业务幂等

Kafka Producer 侧则通过:

1
2
3
PID
+
Sequence Number

避免重试导致消息重复写入。

Kafka Exactly Once 再进一步通过:

1
2
3
4
5
幂等 Producer
+
Transaction
+
Consumer Offset 与输出结果原子提交

保证 Kafka 内部流处理链路中的 Exactly Once。

Consumer 幂等解决的是“同一条消息被重复处理怎么办”。

Producer 幂等解决的是“同一次发送因为重试被 Broker 写入多次怎么办”。

而 Exactly Once,则是在这两类问题之上进一步解决:

消费进度和处理结果如何保证原子一致。


消息队列中的幂等:从重复消费到 Kafka Exactly Once
http://example.com/2026/09/15/消息队列中的幂等:从重复消费到-Kafka-Exactly-Once/
作者
lorixyu
发布于
2026年9月15日
许可协议