☰
CAP 幂等性深入解析:At-Least-Once 投递保证与消费者幂等设计实践
2026/10/9 10:13:11 网站建设 项目流程
  • 后端
  • 消息队列
  • 微服务

【免费下载链接】CAP

基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。

项目地址:https://gitcode.com/dotnetcore/CAP
点击查看免费下载

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),并通过有限次重试驱动最终一致性;但受限于工作事务的概念本质、成功消息的定时清理以及业界惯例,框架不内置严格幂等。因此,生产实践中的正确姿势是:

  1. 明确 CAP 的重试参数(FailedRetryCount、FailedRetryInterval、FailedThresholdCallback、SucceedMessageExpiredAfter),并在 CapOptions.cs 中按业务调整;
  2. 为订阅者设计幂等策略:业务可幂等者优先采用自然幂等,业务不可重复执行者必须引入IMessageTracker之类的显式跟踪;
  3. 让"业务写入 + 幂等标记"处于同一事务,从根本上消除重复执行带来的状态偏移。

幂等不是框架替你做的一件事,而是你在设计消费者时必须承担的职责——理解 At Least Once 的投递模型,是迈出正确设计的第一步。关于 CAP 的事务与 Outbox 写入侧的更多细节,可进一步阅读 transactions.md;订阅者消息映射与头部信息可参考 messaging.md。

  • 后端
  • 消息队列
  • 微服务

【免费下载链接】CAP

基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。

项目地址:https://gitcode.com/dotnetcore/CAP
点击查看免费下载

相关推荐

上一篇:使用 @midwayjs/one-shot 在 Midway 中执行一次性脚本任务
下一篇:CANN ops-nn ForeachSqrt 算子详解:张量列表逐元素平方根计算的实现与 aclnn 调用指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询