重复消费与幂等:重试、确认丢失与 Rebalance
从一次退款消息被执行两次出发,解释发送重试、确认窗口和消费者交接为何产生重复,展开业务键、唯一约束、事务去重、处理中恢复和外部副作用的幂等边界。
退款消费者处理了一条消息,数据库已经记录退款成功,但确认消息没有到达 Broker。消息系统只知道这次消费没有确认,重新投递是合理的;业务已经退款,第二次退款又是错误的。两边掌握的事实不同,重复就出现在这个边界上。
可靠投递经常允许同一事件再次出现。消费者的目标可以是多次接收、多次尝试,但同一业务动作只产生一次有效结果。幂等设计需要把“同一动作”的身份、持久结果和并发裁决写进业务协议,不能只在回调开头加一个缓存判断。
本文聚焦消息消费,沿退款示例解释重复如何发生、如何处理。通用请求超时与 Exactly-once 的概念可先看 超时、重试与幂等。RocketMQ、RabbitMQ 与 Kafka 的行为依据官方文档;退款表、SQL 和恢复流程是通用设计示例,不是线上实习系统的实现披露。
一、重复发生在发送端,也发生在消费端
发送方发送成功,Broker 已接收,返回确认时网络断开。发送方没有收到确认,继续重发。这会产生两条传输记录,对应同一业务事件。Outbox 发送器也有类似窗口:消息发出后、标记已发送之前宕机,重启后再次发送。RabbitMQ 可靠性指南明确提醒未收到确认的重传可能形成重复。
消费端则是在业务提交后、消费确认前故障。SDK 重试、Broker 重新交付或 Kafka 从已提交位点恢复,都可能再次把已经产生过效果的记录交给业务。若业务提交尚未成功,同一次重投又是必要恢复,不能因为“见过消息”就丢弃。
还有部分成功:一个批次里三条消息,前两条已经提交,第三条失败,批次重试可能再次包含前两条。消费者把消息交给线程池后回调超时,后台任务仍然执行;另一份重投同时开始,两份工作还可能并发生效。
消费者重平衡,也就是 Rebalance,会重新分配分区或队列。旧消费者处理了数据但未完成进度提交,新消费者从持久进度接手,重复处理这段数据很正常。分配已经交接,不代表旧线程的数据库事务和外部请求被撤销。
因此幂等不能只针对“重投标记为真”的消息。生产端重新发送的重复业务事件,可能是一条全新的消息,传输层看起来是首次交付;历史重放、故障修复和人工补发也可能换了消息身份。
二、先把事件身份与业务动作身份分开
消息 ID 标识传输记录,eventId 标识一次业务事实,refundId 标识一次退款动作。三者可以有关联,但不必相同。同一个 eventId 被重新发布时应保持稳定;一笔退款的通知事件、状态更新事件和对账事件,可能有不同 eventId,却都指向同一个 refundId。
只用 Broker 消息 ID 去重,可能挡住同一条记录的重投,却挡不住应用补发的新记录。RocketMQ 的消费幂等建议也区分消息身份与业务字段,并指出去重判断需要考虑原子性。
用订单号作为退款唯一键也可能过粗。一张订单可能发生部分退款、多次售后,不同退款不能互相吞掉。唯一身份应来自业务动作的定义,例如租户、退款操作类型和 refundId;退款再次查询、补发和重试复用相同身份,新增另一笔合法退款则产生新身份。
传输身份:brokerMessageId = 本次传输记录
事件身份:eventId = 退款请求事件 r-event-17
动作身份:(tenantId, operationType=REFUND, refundId=r-17)
参数摘要:orderId、金额、币种等不可变业务参数的规范化摘要
同键不同参数应进入冲突处理。第一次退款十元,第二次沿用 refundId 却请求一百元,不能直接当重复成功,也不能覆盖原参数。摘要的用途是识别这个错误,而不是把每次不同序列化格式都当成新动作;字段顺序、可变追踪信息和无关时间戳不应进入身份判断。
三、数据库唯一约束怎样挡住并发重复
“先查询没处理过,再执行业务”有竞争窗口。两个消费者同时查询,都发现不存在,然后各自扣减余额。查询与写入之间需要一个原子胜者,可以借助数据库唯一约束或受条件保护的状态转移。
对于全部效果都在一个数据库内的任务,可以让去重记录和业务写入处在同一事务中。先插入 (consumerName, eventId) 唯一记录;插入成功后执行业务,最后一起提交。并发冲突方等待原事务结束,再根据已提交结果判定是否完成。
BEGIN
INSERT processed_event(consumer_name, event_id, payload_digest)
修改退款、账务或业务状态
保存本次稳定处理结果
COMMIT
提交消费确认或位点
遇到唯一冲突:回滚当前事务,查询已提交记录
同参数且已完成:返回既有结果,确认消费
参数冲突或状态未完成:进入对应恢复,不冒充成功
唯一冲突后怎样继续事务,取决于数据库和驱动。不要把一个数据库的异常处理方式直接套到另一个数据库;示例采用回滚后查询,避免在已经被标记失败的事务中继续操作。数据库提交结果未知时,也应查询同一键确认,而不是立刻生成新的退款动作。
如果去重记录先独立提交,再更新业务,更新前宕机会留下“已处理但未退款”;反过来先退款,再单独写去重,宕机后会重复退款。两个独立提交无法靠调整顺序消除窗口,原子边界必须覆盖相应效果。
记录中加入 consumerName,是为了让库存、退款和通知三个独立订阅者各自处理同一事件。全局只按 eventId 保存一条“已消费”,可能导致第一个订阅者处理后,其他合法订阅者都被误跳过。事件去重作用域与业务动作唯一约束应分别设计。
四、天然幂等的更新也有前置条件
把订单状态设置成 PAID,看上去重复设置不会改变结果。但它可能把后续 CLOSED 状态覆盖回 PAID,也可能每次都触发积分增加或发券。幂等需要审查整段业务效果,不只看最后一条赋值语句。
对于状态机,可以使用合法前置状态进行条件更新。例如 UNPAID 到 PAID 只允许成功一次,成功方才写出发券 Outbox。后续重复看到 PAID,读取已有支付事实;看到 CLOSED,需要处理迟到支付,而不是无条件覆盖。
对于余额加减、库存扣减这样的增量,仅把操作写成 amount = amount + delta 不幂等。可以为账务流水设置业务动作唯一约束,在同一事务里新增流水和改变余额。流水既提供去重证据,也提供对账依据;缓存记录丢失时不能让同一笔动作重新入账。
对于完整快照投影,可以只接受较新版本,防止重复旧快照导致倒退。但收到版本 9 不代表已经执行版本 8 中的发券副作用。快照覆盖和事件动作是不同契约,具体边界见 顺序消费。
查询接口通常容易重复调用,但也要确认它是否隐含生成资源、刷新计费或改变访问次数。如果所谓“只读查询”会在首次调用时创建外部订单,它实际上也是写操作。幂等应根据副作用判断,不能由接口名称判断。
五、Redis 标记和分布式锁为什么只能覆盖部分问题
一种常见写法是在 Redis 用 SET key value NX EX ttl 标记,再执行业务。标记成功后进程崩溃,后续消费者看到 key 就返回成功,业务被跳过。若先执行业务再写 Redis,业务成功后进程崩溃,后续仍会重复执行。
这不是 Redis 不支持原子命令的问题。Redis 的原子操作覆盖 Redis 内部,不能自动覆盖另一个数据库或外部付款渠道。Lua 可以原子完成多个 Redis 步骤,但不能让已经发出的 HTTP 退款成为同一个事务的一部分。
锁可以减少同时执行,但锁过期不能证明旧执行者已经停止。旧执行者因为长时间 GC 暂停,租约到期后新执行者接管,旧执行者恢复后仍可能继续扣款。即使释放锁检查了 owner,已经产生的业务副作用也不会撤回。
可以把 Redis 作为快速筛选或并发削峰,持久业务记录负责最终判定。缓存命中应能对应可查询的已完成事实;缓存缺失只意味着需要访问权威记录,不能意味着“肯定没执行过”。对于必须限制旧执行者的本地写入,可以把领取世代写进条件更新,由持久层拒绝过期世代。
幂等键的保留期也很关键。消息可以积压数天,死信可以数周后重放,缓存标记如果十五分钟就过期,旧动作重新出现时保护已不存在。保留期应覆盖允许的重放范围;重要业务可以把唯一性绑定到退款流水生命周期,而不依赖短期缓存 TTL。
六、外部退款需要一个可查询的持久状态机
调用支付渠道不能与本地数据库自然组成一个事务。渠道已经退款,本地记录尚未提交时宕机,恢复后不知道是否退款;先记录 SUCCESS 再调用,则可能永远没有真正退款。处理这个边界,需要渠道支持的稳定请求身份和结果查询。
一种设计是先持久保存退款动作及不可变参数,状态为 PENDING;领取后转为 PROCESSING,携带相同 refundId 调用渠道。拿到确定成功后保存渠道流水与本地结果;拿到确定拒绝后保存终态错误;请求超时转为 UNKNOWN,按同一 refundId 查询渠道状态。
如果渠道承诺同一键的重复请求返回同一笔动作,就可以在约定窗口和参数约束下安全重试。但要确认幂等键是否真的覆盖退款、有效期多久、是否要求金额完全一致、退款查询是否最终可见。一个接口接受 requestId,不等于它提供业务幂等保证。
处理中状态还需要恢复责任。领取者宕机,后台根据租约和当前证据接管;不能因为 PROCESSING 已存在就永久拒绝,也不能直接把它当成成功。如果外部请求还在进行,新执行者要先查询或使用同一幂等键,避免把接管变成另一笔退款。
渠道不支持幂等也不支持可识别查询时,本地去重不能证明外部只生效一次。可以缩小自动重试范围,保留调用证据,进入对账或人工确认;不能把结果未知直接标成失败,然后无限发起退款。
MQ 确认的责任边界可以是退款真正完成,也可以是持久退款工作流已经可靠接管。后一种更容易避免长回调,但接管记录必须可恢复,有独立调度和告警。消费者返回成功之后,后续完成责任已经从 MQ 转移到业务系统。
七、Kafka 位点和 Exactly-once 不能替代业务幂等
Kafka 消费者读取位置和提交位置不同。poll 返回数据后,读取位置已经推进;提交位置是恢复依据,不表示每条业务都已经落库。4.1 API 还明确说明 offset 可能不连续,因此不能把每一个整数缺口都当作消息丢失。KafkaConsumer
先提交位点再写数据库,故障窗口会漏业务;先写数据库再提交位点,故障窗口会重业务。常见方案是关闭不符合任务语义的自动提交,在本地事务中完成事件去重和业务,再提交相应分区的完成位置。位点提交失败时允许重读,由数据库识别重复。
并行处理同一分区时,要跟踪连续完成范围。offset 100 还没完成,101、102 已完成,不能为了追赶进度直接提交 103。恢复会从 103 开始,跳过尚未完成的 100。这里的连续范围是日志消费进度,不要求实际记录整数完全连续。
Kafka 幂等生产者解决的是其协议范围内重试导致的重复追加,不会识别应用主动发送的两个独立业务请求。Kafka 事务可以把 Kafka 输入位点与 Kafka 输出记录协调提交,消费者还要使用相应隔离级别;外部数据库和渠道并不会自动加入这个事务。Kafka 消息交付语义
如果输出就是 Kafka Topic,事务非常有用;如果输出是退款或数据库流水,需要目的系统配合。不能从“启用 Exactly-once”推导“消费者代码永远只运行一次”,更不能推导“第三方一定只扣一次款”。
八、重平衡与历史重放怎样守住相同动作
消费责任交接时,旧消费者应停止接受新工作,并对已开始任务采取明确策略:等待有界时间、持久接管或让未完成记录重投。即使交接协议做得很好,进度确认窗口仍会产生重复,因此业务幂等必须独立成立。
旧工作者晚到的提交要靠业务条件或执行世代约束。新的处理者完成 refundId=r-17 后,旧处理者不能把它从 SUCCESS 改回 PROCESSING,也不能用旧参数覆盖渠道流水。恢复必须尊重终态以及参数不变式。
人工重放不应该随意改 eventId。若目的是恢复同一原始事件,就保留事件身份与业务动作身份,另外记录 replayBatchId 追踪这次操作;如果为了绕过去重而换新 refundId,会把恢复操作变成新的退款意图。
确实需要重建投影时,可以换独立的消费者作用域或投影世代,重新计算展示数据;退款、发券、短信等副作用则不能跟着一起“重建一次”。投影恢复和真实动作补发应该走不同控制路径。
去重记录归档后,还应定义超出允许范围的旧事件怎么处理。可以查询归档业务流水、拒绝过期动作并人工确认,或使用长期动作唯一约束。只清理去重表,却允许任意历史消息直接进入增量消费者,是一个隐蔽的重复入口。
九、沿一笔退款检查完整性
退款请求 r-17 进入 Outbox,发送确认丢失后重复发布。消费者第一次持久建立退款动作,第二次发现同键同参数,读取已有状态而不新建;执行者用 r-17 请求渠道,渠道退款后响应丢失,本地进入 UNKNOWN。
恢复任务查询到渠道已有退款流水,保存 SUCCESS 和渠道凭证。原消息再投递时,消费者返回既有成功结果并确认。即使重平衡让两个执行者短暂重叠,渠道幂等身份和本地终态条件也限制同一动作再次生效。
需要观察的指标包括事件去重命中、参数冲突、长时间 PROCESSING、UNKNOWN 年龄、同业务键的渠道流水数量和恢复失败。去重命中增加不一定是异常,可能只是故障后正常重放;重复业务效果出现才是明确的不变量破坏。
我会先问一个具体问题:宕机后,下一位执行者能否仅凭持久记录判断同一动作做到哪一步?如果答案仍依赖旧进程内存、短期缓存或“应该成功了”,幂等方案就还没有覆盖恢复。持久身份、原子结果和可查询外部证据,才能让重复投递成为可处理的正常路径。
参考资料
- RocketMQ 基础最佳实践:消息身份与业务幂等建议。
- RabbitMQ 可靠性指南:发送确认丢失与重新交付。
- Kafka 4.1 Consumer API:读取位置、提交位置与恢复。
- Kafka 4.1 设计:幂等与事务的交付语义边界。
如果这篇文章对你有帮助