MQ消息的重复消费与幂等
消息幂等性,避免MQ消息重复消费
摘要:消息队列(MQ)是分布式系统的基石,但“至少一次投递”的语义决定了消息重复是常态,幂等是必须。本文从重复消费的根源出发,深入剖析三种常见幂等策略(消费记录表、Redis+DB、数据库唯一索引)的适用场景与陷阱,并提供一套可落地的选型指南。
一、“至少一次投递”
在生产环境中,我们几乎不可能保证消息“只被消费一次”。更常见的语义是 “至少一次投递(At Least Once)”,即消息可能被重复投递,消费者必须自行保证幂等性——多次消费同一条消息,最终的业务影响与消费一次完全相同。
重复消费的根本原因,在于消费者未能正确地(或及时地)向 Broker 返回确认(ACK)。当 Broker 未收到 ACK 时,它会认为消费者“挂掉了”,从而将消息重新投递给其他消费者或稍后重试。
具体来看,以下三种场景都可能导致 ACK 失败:
| 场景 | 描述 | 典型表现 |
|---|---|---|
| 拿到消息后 | 消息已拉取到客户端,但业务尚未执行,进程崩溃或网络闪断 | 应用 OOM、K8s Pod 重启 |
| 消费中 | 业务逻辑执行过程中抛出异常,未走到 ACK 代码 | DB 超时、远程 RPC 调用失败 |
| 消费完成后 | 业务执行成功,但在回传 ACK 时网络故障 | 客户端与 Broker 连接断开 |
注意:第一种情况属于系统级异常,往往无法被业务代码捕获,通常由 MQ 客户端自身的重试机制兜底。我们设计幂等策略时,主要聚焦于业务执行过程中的错误。
二、三大幂等策略
业界主流方案可归纳为以下三类,它们分别适用于不同的业务场景:
| 策略 | 核心思想 | 适用场景 |
|---|---|---|
| 消费记录表 | 记录每一条消息的处理状态,通过状态流转控制重复请求 | 需要对消息处理状态进行精确追溯的场景 |
| 数据库唯一业务ID | 利用数据库唯一约束或乐观锁版本号,物理层面拦截重复写入 | 数据插入/更新类操作,尤其是核心交易链路 |
| 变更目标状态机 | 基于业务对象的状态流转,只允许特定状态下的操作 | 订单、工单等有明确生命周期管理的场景 |
下面逐一深入分析,并重点揭示其中的陷阱与应对方案。
三、策略一:消息消费记录表
3.1 基础实现:状态表 + 状态流转
创建一个消息消费记录表,以消息唯一 ID(或业务流水号)作为主键,记录每次消费的状态:
message_consume_record
├── msg_id (PK)
├── status (PROCESSING / COMPLETE / FAILED)
├── create_time
└── update_time核心流程:
- 收到消息,插入一条
status = PROCESSING的记录。 - 执行业务逻辑。
- 成功则更新状态为
COMPLETE,失败则更新为FAILED(或直接删除记录,让重试时重新插入)。 - 重复消息进来时,查询记录表:
- 若状态为
COMPLETE→ 直接 ACK,视为重复。 - 若状态为
PROCESSING→ 阻塞等待或拒绝重试,等待第一条处理完成。
- 若状态为
⚠️ 并发处理的关键:当重复消息查询到
PROCESSING状态时,绝不能直接返回成功,否则第一条消息如果后续失败,数据将丢失。正确的做法是:
- 方案A:自旋等待 + 短暂 sleep,轮询直到状态变更为
COMPLETE或FAILED。- 方案B:直接抛出异常,让 MQ 稍后重试(利用退避机制)。
3.2 高并发优化:Redis 缓存前置
当 QPS 较高时,每次查询数据库会带来较大压力。因此引入 Redis 作为前置去重缓存:
- 收到消息后,先执行
SETNX msg_id(带上合理的 TTL,如 10 分钟)。 - 若
SETNX成功,则继续执行业务。 - 若
SETNX失败,直接 ACK,视为重复。
致命陷阱(务必警惕):
1. 消费者收到消息
2. Redis SETNX 成功 ✅
3. 执行 DB 事务 → 超时/回滚 ❌
4. 消费者结束,Redis Key 仍然存在
5. MQ 重试 → 再次 SETNX 失败 → 直接返回成功 ❌
6. 最终:数据库无数据,消息被丢弃 → 数据丢失!正确使用 Redis 缓存的三大铁律:
| 原则 | 具体做法 |
|---|---|
| 数据库兜底 | 数据库中必须保留唯一索引(或业务主键),Redis 只为拦截 99% 的流量,最终一致性必须由 DB 保证 |
| 异常删除 | 在 catch 代码块中,务必显式删除 Redis Key,让重试消息能够再次执行 |
| 短 TTL | 设置合理的过期时间(如 10 分钟),作为最后一道容错防线 |
🚫 红线原则:对于高价值业务(如支付、订单扣减),绝不能仅依赖 Redis 做幂等判断。Redis 只能作为缓存加速,最终决策必须回归数据库。
四、策略二:数据库唯一业务ID(必要)
这是成本最低、可靠性最高的幂等策略,其核心在于利用数据库的唯一约束或乐观锁版本号,从物理层面拦截重复写入。
4.1 插入场景:唯一索引
在业务表中,为**业务流水号(BizId)**建立唯一索引:
CREATE TABLE order_tbl (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
order_no VARCHAR(32) UNIQUE NOT NULL, -- 唯一索引
amount DECIMAL(10,2),
status TINYINT
);- 第一条消息插入成功。
- 重复消息再次插入时,数据库抛出
DuplicateKeyException,消费者捕获后直接 ACK。 - 无需任何额外查询,一次写入即可完成幂等判断,性能最优。
4.2 更新场景:乐观锁版本号
对于数据更新操作,使用 版本号(Version) 作为乐观锁:
UPDATE account_tbl
SET balance = balance - 100, version = version + 1
WHERE biz_id = 'ORDER_20250816_001' AND version = 5;- 第一次更新:
version从 5 → 6,影响行数 = 1。 - 重复消息:
version仍为 5,但数据库当前版本已是 6,影响行数 = 0 → 判定为重复请求,直接 ACK。
✅ 结论:数据库唯一约束/乐观锁是幂等设计的终极兜底,它不依赖任何外部组件,且能保证数据的强一致性。在生产环境中,即使使用了 Redis 或消费记录表,也应当保留数据库层的唯一约束作为最后一道防线。
五、策略三:变更目标状态机
当业务对象具有明确的状态生命周期时,可以利用状态机天然实现幂等。
典型案例:订单状态流转 待支付(PENDING) → 已支付(PAID) → 已发货(SHIPPED)。
- 支付消息仅当订单状态为
PENDING时才能执行扣款并变更为PAID。 - 若订单已经是
PAID,说明该支付消息已经处理过了,直接 ACK。 - 实现方式:SQL 条件更新,天然保证幂等性。
UPDATE order_tbl
SET status = 'PAID', pay_time = NOW()
WHERE order_no = 'O20250816001' AND status = 'PENDING';- 第一次执行:影响 1 行,成功。
- 重复执行:
status已不是PENDING,影响 0 行,判定为重复,直接 ACK。
这种方案本质上是业务层的 CAS(Compare-And-Swap),无需额外存储,且语义清晰。适用于订单、工单、审批流等所有有状态管理的场景。
六、策略对比与选型指南
| 策略 | 性能开销 | 可靠性 | 复杂度 | 推荐场景 |
|---|---|---|---|---|
| 消费记录表 | 中(两次 DB 操作) | 高 | 中 | 需要精确追踪消息处理状态的场景 |
| Redis + DB | 低(一次 Redis + 一次 DB) | 高(需异常删除 + TTL) | 高 | 高并发、可接受少量代码复杂度的场景 |
| 数据库唯一索引 | 低(一次 DB 写入) | 极高(物理保障) | 低 | 所有核心业务场景的必选项 |
| 乐观锁状态机 | 低(一次 DB 条件更新) | 高 | 低 | 有明确状态流转的业务(订单/工单) |
实战黄金法则:
- 核心交易链路:必须使用数据库唯一索引 + 乐观锁作为最终兜底。
- 高并发场景:可在 DB 之上叠加 Redis 缓存,但务必遵守“异常删除 + 短 TTL + 数据库兜底”三原则。
- 有状态业务:优先使用状态机,这是天然幂等,无需额外设计。
- 需要精确追溯:采用消费记录表 + PROCESSING 状态阻塞等待,但要注意性能损耗。
七、FAQ:批量消费场景下的幂等如何处理?
批量拉取(如一次拉取 100 条)时,如果其中第 50 条处理失败,前 49 条已成功,应该怎么办?
- ❌ 错误做法:直接对整批 ACK → 失败的那条就丢了;或整批 NACK → 前 49 条会重复消费。
- ✅ 正确做法:将批量拆解为单条处理 + 逐条确认,或利用 MQ 的死信队列,只将失败的消息发回重试,成功消息正常提交偏移量。
这也是 RocketMQ 中 ConsumeOrderly(顺序消费)与 ConsumeConcurrently(并发消费)模式区别的核心考量之一。
八、总结
消息幂等不是一个可选的“优化”,而是基于 MQ “至少一次投递”语义下的强制性要求。本文从重复消费的三大根源出发,对比分析了四种主流幂等策略,并重点强调了:
- Redis 只能做缓存,不能做最终决策。
- 数据库唯一索引 / 乐观锁是最可靠的兜底方案。
- 状态机是天然幂等,应优先使用。
希望这份总结能帮助你在实际项目中设计出既高效又可靠的消息幂等方案。如果文章对你有帮助,欢迎转发给需要的朋友,也欢迎在评论区交流你在实践中遇到的幂等难题。
关于作者:多年分布式系统开发经验,专注于消息中间件、微服务架构和云原生技术。本文基于实际生产环境的踩坑经历总结而成。
本文首发于个人技术博客,转载请注明出处。