
后端消息队列微服务【免费下载链接】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 * 360086400 秒即 24 小时失败消息的默认过期时间FailedMessageExpiredAfter为15 * 24 * 360015 天二者均可通过配置调整。换句话说成功消费的记录只保留有限时间这一设计至今成立具体保留时长取决于你的配置。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把状态更新为FailedUpdateMessageForRetry中重试阈值取Math.Min(_options.FailedRetryCount, 3)——FailedRetryCount默认 50 次但前几次失败会立即触发快速重试当重试次数达到FailedRetryCount时会触发FailedThresholdCallbackFailedInfo回调此时消息被判定为永久失败、不再重试特殊情况下若异常为SubscriberNotFoundException找不到订阅者消息会直接放弃重试。发布侧的逻辑与之对称见 IMessageSender.Default.csSetSuccessfulState设置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 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐CAP 消息幂等性深度解析交付语义、重试机制与消费端幂等设计实践CAP 消息幂等性深度解析交付语义、重试机制与消费端幂等设计实践 CAPDotNetCore.CAP是一套基于最终一致性的分布式事务解决方案同时也内置了后端消息队列微服务消息路由CAP 分布式事件总线中的消息幂等性交付保证、重试机制与消费端幂等实践指南CAP 分布式事件总线中的消息幂等性交付保证、重试机制与消费端幂等实践指南 在基于最终一致性的事件总线 CAP基于 Outbox 模式中消费端收到同一条后端消息队列微服务终极指南Disque分布式消息投递语义详解——At-Least-Once与At-Most-Once实现原理终极指南Disque分布式消息投递语义详解——At Least Once与At Most Once实现原理 Disque作为一款高性能分布式消息代理其核心价消息队列后端上一篇使用 midwayjs/one-shot 在 Midway 中执行一次性脚本任务下一篇CANN ops-nn ForeachSqrt 算子详解张量列表逐元素平方根计算的实现与 aclnn 调用指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考