- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
CAP 是一款基于最终一致性理念的分布式事务解决方案与 Outbox 模式事件总线,其核心文档 idempotence.md 系统阐述了框架的投递保证模型(At Least Once)以及为何不内置严格幂等,并给出了两种实用的消费者幂等设计路径。本文以该文档为骨架,结合仓库源码(CapOptions.cs、IProcessor.NeedRetry.cs、ISubscribeExector.Default.cs 等)逐层展开,帮助你彻底理解 CAP 的重试与消息生命周期机制,并在自己的订阅者(Consumer)中落地可靠的幂等处理方案。
交付保证(Delivery Guarantees):先理解消费端会收到几次消息
在讨论幂等性之前,必须先把消费端的消息交付模型讲清楚。CAP 不使用 MS DTC 或任何形式的 2PC(两阶段提交)分布式事务机制,因此存在"消息至少被严格交付一次"这一固有局限。基于消息的系统中,交付保证通常存在以下三种可能:
- Exactly Once(仅有一次)——带
*号,因为在通用场景下它根本无法实现 - At Most Once(最多一次)
- At Least Once(最少一次)
At Most Once:最多一次
"最多一次"保证你要么收到全部消息,要么一条都收不到,不会重复。
这种保证可能来自消息系统与你的代码按如下顺序执行:
1. 从队列移除消息 2. 开始一个工作事务 3. 处理消息(你的代码) 4. 是否成功? Yes: 1. 提交工作事务 No: 1. 回滚工作事务 2. 将消息放回队列在理想情况下,这套流程运转得很好——消息被接收、工作事务被提交、一切皆大欢喜。
但现实往往不会如此顺利,尤其是当你处理大量工作时。考虑一下:如果在步骤 1 之后发生任何故障,当你试图执行步骤 4/2(把消息放回队列)时——网络临时不可用、消息代理(Broker)重启,或者主机因为系统更新而重启——消息就彻底丢失了。
如果这正是你想要的,那也无妨;但 CAP 中绝大多数概念都围绕**持久消息(DURABLE messages)**展开——这些消息的内容重要程度与数据库中的数据相当。
At Least Once:最少一次
"最少一次"保证一旦出现故障,你会收到全部消息一次或多次。
这需要略微调整执行顺序,并且要求消息队列系统支持事务或 ACK 机制:要么是传统的 begin-commit-rollback 协议(MSMQ 如此),要么是 receive-ack-nack 协议(RabbitMQ、Azure Service Bus 等如此)。大致流程如下:
1. 抢占队列中的消息(获取 lease) 2. 开始一个工作事务 3. 处理消息(你的代码) 4. 是否成功? Yes: 1. 提交工作事务 2. 从队列删除消息 No: 1. 回滚工作事务 2. 释放队列中消息的抢占只要步骤 1 中"抢占(lease)"带有恰当的过期时间,那么无论情况多么糟糕,我们都能够保证:只有当"工作事务"成功提交后,消息才会真正从队列中删除(步骤 4/2)。失败或抢占超时时,消息总能被再次接收,从而确保工作事务最终提交成功。
什么是"工作事务"?
"工作事务"并不特指关系型数据库事务,它是一个概念——代表执行代码的原子性。它可能是:
- 传统关系型数据库中的事务(对这一场景的支持历来很好);
- 支持事务的文档数据库中的事务(如 RavenDB、PostgreSQL、MongoDB);
- 一个概念性事务,代表你处理消息所产生的后果:更新 MongoDB 中的文档、移动文件系统中的文件、修改内存中的数据结构等。
正是因为"工作事务"是一个概念性事实,"工作事务"与"队列事务"(与消息队列系统之间的协议)无法原子化地同时提交或回滚,才导致Exactly Once 在通用场景下不可能实现——不存在某种机制能把二者原子化地保持一致。
CAP 中的幂等性:为什么框架不内置严格幂等
CAP 采用的交付保证是At Least Once(相关表述可见 英文文档 与 中文文档)。
由于 CAP 拥有临时存储介质(数据库表),理论上可以实现 At Most Once;但为了严格保证消息不丢失,CAP 没有提供相关功能或配置。文档从以下四个层面解释了为什么 CAP 没有实现(或达成)幂等:
1. 消息写入成功,但执行 Consumer 方法失败
Consumer 方法执行失败的原因非常多。如果不知道具体场景,盲目地选择重试或不重试都是不正确的。
典型例子:假如消费者是一个扣款服务,扣款已经成功执行,但写扣款日志时失败了。此时 CAP 会判定为"消费者执行失败"并进行重试。如果客户端自己没有保证幂等性,框架的重试必然造成多次扣款的严重后果。
2. Consumer 方法执行成功,但又收到了同一条消息
这个场景同样存在:Consumer 最初已经执行成功,但由于某种原因(如 Broker 宕机恢复),相同的消息又被重新接收。CAP 收到 Broker 消息后会将其视为一条新消息,再次对 Consumer 执行。因为它是新消息,此时 CAP 同样无法做到幂等。
3. 当前的数据存储模式无法做到幂等
CAP 存储消息的表中,成功消费的消息会在一定时间后被清理(文档描述为约 1 小时后删除),因此对于历史性消息无法进行幂等校验。所谓历史性消息,是指 Broker 由于某种原因维护、或人工处理过的消息——此时无法验证它们是否已被处理过。
当前仓库的实际情况:在 CapOptions.cs 中,成功消息的默认过期时间
SucceedMessageExpiredAfter为24 * 3600(86400 秒,即 24 小时),失败消息的默认过期时间FailedMessageExpiredAfter为15 * 24 * 3600(15 天),二者均可通过配置调整。换句话说,"成功消费的记录只保留有限时间"这一设计至今成立,具体保留时长取决于你的配置。
4. 业界做法
许多基于事件驱动的框架都要求用户自己保证幂等性操作,例如 ENode、RocketMQ 等。
结论是:从实现角度来说,CAP 可以提供一些不那么严格的幂等,但严格幂等无法做到。
源码视角:CAP 的重试机制与消息状态机
文档所述"执行失败会重试"并非泛泛而谈,仓库源码中有完整的落点。理解这些机制,有助于你判断为什么必须在业务侧做幂等。
失败消息的持久化与重试处理器
订阅者执行失败后,状态会被标记为Failed,并进入重试队列。核心处理器是 MessageNeedToRetryProcessor:
- 它以
FailedRetryInterval(默认 60 秒)为轮询间隔,从存储中取出需要重试的消息; - 发布侧消息通过
_dispatcher.EnqueueToPublish(message)重新投递,消费侧消息通过_dispatcher.EnqueueToExecute(message)重新执行; - 多实例部署时,可通过
UseStorageLock开启分布式存储锁,确保集群中只有一个实例执行重试,避免重复处理(这正是 CapOptions.cs 中UseStorageLock的用途)。
重试次数、阈值与兜底回调
在 ISubscribeExector.Default.cs 中可以看到消费侧失败处理的完整逻辑:
- 每次失败调用
SetFailedState,将message.Retries递增,并通过ChangeReceiveStateAsync把状态更新为Failed; UpdateMessageForRetry中,重试阈值取Math.Min(_options.FailedRetryCount, 3)——FailedRetryCount默认 50 次,但前几次失败会立即触发快速重试;- 当重试次数达到
FailedRetryCount时,会触发FailedThresholdCallback(FailedInfo回调),此时消息被判定为永久失败、不再重试; - 特殊情况下,若异常为
SubscriberNotFoundException(找不到订阅者),消息会直接放弃重试。
发布侧的逻辑与之对称,见 IMessageSender.Default.cs:SetSuccessfulState设置Succeeded状态与过期时间,SetFailedState则累计重试次数并写入失败状态。
消息清理器:成功消息为何只保留有限时间
CollectorProcessor 负责清理过期数据:它以CollectorCleaningInterval(默认 300 秒)为周期,通过IDataStorage.DeleteExpiresAsync分批(每批 1000 条)删除Published与Received表中过期的消息记录(表名由 IStorageInitializer 提供)。这正是文档第 3 点所述"成功消费的消息会在一定时间后被删除"的底层实现——也是 CAP 无法对历史消息做幂等校验的直接原因。
综合以上源码事实,可以更清晰地理解 CAP 的完整状态机:发布/接收 → 重试(有限次数)→ 成功(限时保留)或永久失败(触发回调)。框架保证"不丢失",但不保证"不重复",重复的兜底责任必须由业务侧承担。
方案一:以自然的方式处理幂等消息
通常情况下,让"消息被执行多次而不会产生意外结果"最自然的方式,是采用操作对象自带的幂等功能。例如,处理一条消息本质上就是调用领域对象上的幂等方法:
obj.MarkAsDeleted();或
obj.UpdatePeriod(message.NewPeriod);利用数据库提供的INSERT ON DUPLICATE KEY UPDATE(或 PostgreSQL 的ON CONFLICT、SQL Server 的MERGE等)可以很轻松地达成这种效果:无论消息被消费多少次,第二次插入命中主键/唯一键时只做更新或直接忽略,业务状态不会发生偏移。这种方式不需要额外的状态存储,实现成本最低,适合业务本身可天然幂等的场景。
方案二:显式处理重复投递(IMessageTracker 模式)
另一种让消息处理具备幂等性的方式,是显式跟踪已处理消息的 ID,然后在代码中处理重复投递。
基本思路是:在消息传递过程中带上唯一 ID,由独立的"消息跟踪器"记录每个消息 ID 的处理状态。假设你使用与业务工作共享同一事务数据存储的IMessageTracker,代码大致如下:
readonly IMessageTracker _messageTracker; public SomeMessageHandler(IMessageTracker messageTracker) { _messageTracker = messageTracker; } [CapSubscribe] public async Task Handle(SomeMessage message) { if (await _messageTracker.HasProcessed(message.Id)) { return; } // 在这里执行实际工作 // ... // 记录该消息已处理 await _messageTracker.MarkAsProcessed(message.Id); }要点拆解:
- 判断先行:进入订阅方法后首先调用
HasProcessed(message.Id),已处理则直接返回,避免重复执行业务逻辑; - 事务一致是关键:
IMessageTracker的记录操作应与业务操作处于同一事务中——如果业务提交成功而记录未提交,重复投递仍会触发二次执行;反之亦然。这也是文档强调"使用与其余工作相同的事务数据存储"的原因; - 落库收尾:业务执行完成后调用
MarkAsProcessed写入处理状态。
对于IMessageTracker的具体实现,可以使用 Redis、数据库等存储消息 ID 及对应的处理状态(如:以消息 ID 为 key 的SETNX,或数据库中的唯一约束表)。唯一约束表配合事务写入,可以天然防止并发下的重复处理。
方案对比与选型建议
| 方案 | 实现成本 | 适用场景 | 关键前提 |
|---|---|---|---|
自然幂等(幂等方法 /INSERT ON DUPLICATE KEY UPDATE) | 低 | 业务操作本身可重复执行且结果一致(如标记删除、幂等更新) | 领域方法或数据库语句真正幂等 |
显式跟踪(IMessageTracker) | 中 | 业务操作不可重复(如扣款、发券、转账) | 跟踪记录与业务写入处于同一事务/原子操作 |
总结:在 CAP 中正确面对重复消息
CAP 基于 Outbox 模式与数据库事务保证消息不丢失(At Least Once),并通过有限次重试驱动最终一致性;但受限于工作事务的概念本质、成功消息的定时清理以及业界惯例,框架不内置严格幂等。因此,生产实践中的正确姿势是:
- 明确 CAP 的重试参数(
FailedRetryCount、FailedRetryInterval、FailedThresholdCallback、SucceedMessageExpiredAfter),并在 CapOptions.cs 中按业务调整; - 为订阅者设计幂等策略:业务可幂等者优先采用自然幂等,业务不可重复执行者必须引入
IMessageTracker之类的显式跟踪; - 让"业务写入 + 幂等标记"处于同一事务,从根本上消除重复执行带来的状态偏移。
幂等不是框架替你做的一件事,而是你在设计消费者时必须承担的职责——理解 At Least Once 的投递模型,是迈出正确设计的第一步。关于 CAP 的事务与 Outbox 写入侧的更多细节,可进一步阅读 transactions.md;订阅者消息映射与头部信息可参考 messaging.md。
- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
相关推荐
CAP 消息幂等性深度解析:交付语义、重试机制与消费端幂等设计实践
CAP 消息幂等性深度解析:交付语义、重试机制与消费端幂等设计实践 CAP(DotNetCore.CAP)是一套基于最终一致性的分布式事务解决方案,同时也内置了
后端消息队列微服务消息路由CAP 分布式事件总线中的消息幂等性:交付保证、重试机制与消费端幂等实践指南
CAP 分布式事件总线中的消息幂等性:交付保证、重试机制与消费端幂等实践指南 在基于最终一致性的事件总线 CAP(基于 Outbox 模式)中,消费端收到同一条
后端消息队列微服务终极指南:Disque分布式消息投递语义详解——At-Least-Once与At-Most-Once实现原理
终极指南:Disque分布式消息投递语义详解——At Least Once与At Most Once实现原理 Disque作为一款高性能分布式消息代理,其核心价
消息队列后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考