- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
CAP(dotnetcore/CAP)是一个开箱即用的 .NET 事件总线(EventBus)框架,同时为微服务或 SOA 系统提供基于最终一致性的分布式事务解决方案。它在微软 eShop 展开,结合仓库源码与示例,系统讲解 CAP 的核心设计理念、工作原理、可扩展架构与实战用法,帮助读者理解"为什么直接用消息队列不够可靠,而 CAP 可以",并掌握发布订阅、事务集成、配置调优的完整方法。
CAP 是什么:EventBus 与分布式事务的一体化框架
CAP 定位于两个核心场景:
- EventBus:在微服务、SOA 系统中解耦服务间的异步事件通信;
- 分布式事务框架:基于最终一致性(Eventual Consistency)模型,解决跨服务数据一致性问题。
它采用Outbox 模式(本地消息表)实现事务可靠性。核心思想是:业务数据与待发送事件消息在同一个数据库事务内写入本地消息表,消息由后台处理器可靠地投递到消息队列;下游消费者消费成功后确认,未成功则自动重试。这一设计保证了"事件消息永不丢失",从而规避了单纯使用消息队列时"先发消息再写库、或先写库再发消息"所带来的双写不一致问题。
服务A业务操作 ──同库事务──> 本地消息表(Outbox) │ 后台投递 ▼ 消息队列(Transport) │ 消费 ▼ 服务B订阅处理 ──确认──> 本地接收表 ──重试/补偿──> 最终一致什么是 EventBus:组件解耦的通信机制
介绍文档对 EventBus 的定义非常精炼:事件总线是一种允许不同组件彼此通信而无需彼此了解的机制。
- 发布者向总线发送事件,无需知道"谁在接收、有多少人接收";
- 订阅者监听总线上的事件,无需知道"事件由谁发出";
- 组件之间因此可以相互通信而不产生强依赖;
- 替换某个组件时,只要新组件了解正在收发的事件格式,其他组件完全无感知。
这一机制正是微服务系统"可扩展、可靠、易于更改"的基石:服务边界清晰、故障隔离、演进成本低。
CAP 的差异化特色:约定大于配置,零侵入式使用
与许多 Service Bus 或 Event Bus 框架不同,CAP 有两点突出特色:
- 不需要继承或实现任何接口即可收发消息。发布方只需注入
ICapPublisher调用发布方法;订阅方只需在方法上标注[CapSubscribe]特性。这带来极高的灵活性。 - 约定大于配置。默认值经过精心设计(见下文配置表),大多数场景开箱即用,对新手极其友好,同时保持轻量级。
唯一的"例外"是:非 Controller 的业务类订阅者需要实现标记接口ICapSubscribe(仅作标记、无任何成员),用于应用启动时自动发现订阅类,参见 ICapSubscribe.cs。Controller 中的 Action 则完全无需实现任何接口。
模块化架构:可插拔的传输、存储与序列化
CAP 采用模块化设计,具有高度的可扩展性:消息队列、存储、序列化方式等系统元素均可替换为自定义实现。这一设计在仓库src/目录结构中得到直接印证——每个功能点对应独立项目:
| 维度 | 抽象接口(位于 src/DotNetCore.CAP) | 内置实现 |
|---|---|---|
| 消息队列(Transport) | ITransport、IConsumerClient、IConsumerClientFactory(Transport 目录) | Kafka、RabbitMQ、Azure Service Bus、AmazonSQS、NATS、Pulsar、Redis Streams、InMemoryMessageQueue |
| 存储(Storage) | IDataStorage、IStorageInitializer(Persistence 目录) | SqlServer、MySql、PostgreSql、MongoDB、InMemoryStorage |
| 序列化 | ISerializer(Serialization 目录) | 默认基于System.Text.Json(JsonSerializerOptions可自定义) |
| 监控 | IMonitoringApi(Monitoring 目录) | Dashboard 实时看板、Consul / K8s 节点发现、OpenTelemetry |
以存储抽象为例,IDataStorage定义了消息落库、状态流转、重试查询、过期清理、分布式锁等完整契约(见 IDataStorage.cs),而IStorageInitializer负责建表与表名管理(GetPublishedTableName/GetReceivedTableName/GetLockTableName,见 IStorageInitializer.cs)。接入一种新数据库或新消息队列,只需实现对应接口并注册扩展(ICapOptionsExtension)即可。
发布消息:ICapPublisher 与事务集成
发布入口是ICapPublisher(见 ICapPublisher.cs),支持同步/异步、带自定义 Header、延迟发布等多种形式:
// 同步发布 capBus.Publish("test.show.time", DateTime.Now); // 异步发布 await capBus.PublishAsync("xxx.services.show.time", DateTime.Now); // 带自定义 Header var header = new Dictionary<string, string>() { ["my.header.first"] = "first", ["my.header.second"] = "second" }; capBus.Publish("test.show.time", DateTime.Now, header); // 延迟发布(不依赖消息队列本身的延迟特性) capBus.PublishDelay(TimeSpan.FromSeconds(100), "test.show.time", DateTime.Now); await capBus.PublishDelayAsync(TimeSpan.FromSeconds(30), "xxx.services.show.time", DateTime.Now);与数据库事务的原子集成(Outbox 核心)
CAP 的核心能力在于将消息发布与业务数据库操作纳入同一个事务。以 ADO.NET 与 EF Core 为例(完整示例见 README.md 与 快速开始):
// ADO.NET 方式 using (var connection = new MySqlConnection(ConnectionString)) using (var transaction = connection.BeginTransaction(_capBus, autoCommit: true)) { // 业务写库…… _capBus.Publish("xxx.services.show.time", DateTime.Now); } // EF Core 方式 using (var trans = dbContext.Database.BeginTransaction(_capBus, autoCommit: true)) { // 业务写库…… _capBus.Publish("xxx.services.show.time", DateTime.Now); }事务封装统一由ICapTransaction抽象提供(见 ICapTransaction.cs):事务提交(Commit/CommitAsync)时缓冲的 CAP 消息才真正发送到消息队列;回滚(Rollback/RollbackAsync)则丢弃全部未提交消息。由此保证消息与业务要么同时成功、要么同时失败,从源头杜绝"双写不一致"。
订阅消息:CapSubscribe 特性与消费模型
订阅方只需在方法上添加[CapSubscribe("topic.name")]特性(定义见 CAP.Attribute.cs),并支持异步签名、CancellationToken、Header 注入:
// Controller 内直接订阅 public class ConsumerController : Controller { [NonAction] [CapSubscribe("test.show.time")] public void ReceiveMessage(DateTime time) { Console.WriteLine("message time is:" + time); } } // 业务服务类需实现 ICapSubscribe public class SubscriberService : ISubscriberService, ICapSubscribe { [CapSubscribe("xxx.services.show.time")] public void CheckReceivedMessage(DateTime datetime) { /* 处理业务 */ } } // 异步 + 取消令牌 + 读取 Header [CapSubscribe("test.show.time")] public async Task ProcessAsync(Message message, [FromCap] CapHeader header, CancellationToken cancellationToken) { Console.WriteLine(header["my.header.first"]); await SomeOperationAsync(message, cancellationToken); }高级订阅特性
从TopicAttribute与CapSubscribeAttribute的实现(TopicAttribute.cs、CAP.Attribute.cs)可以确认以下能力:
- 订阅组
Group:对应 Kafka 的groups.id/ RabbitMQ 的queue.name。默认组名为cap.queue.{入口程序集名小写}(见CapOptions.DefaultGroupName)。同组多实例竞争消费(负载均衡),不同组广播消费(Fan-out); - 部分订阅
isPartial: true:类级主题与方法级主题拼接,如类上[CapSubscribe("customers")]+ 方法上[CapSubscribe("create", isPartial: true)]组合为订阅customers.create; - 并发限制
GroupConcurrent:限制某主题并发消费的消息数,未指定 Group 时会自动以主题名创建组; - 回调订阅
callbackName:发布时可指定回调订阅者名,配合CapHeader.AddResponseHeader/RemoveCallback/RewriteCallback实现请求-响应式交互。
核心配置参数一览(CapOptions)
CapOptions(见 CAP.Options.cs)集中管理处理管线的关键行为,以下为构造函数中确认的默认值:
| 配置项 | 默认值 | 作用 |
|---|---|---|
Version | "v1" | 消息版本标识(≤20 字符),用于多实例/多版本隔离 |
DefaultGroupName | cap.queue.{程序集名小写} | 默认消费者组名 |
SucceedMessageExpiredAfter | 86400(24h) | 成功消息自动清理时间(秒) |
FailedMessageExpiredAfter | 1296000(15 天) | 失败消息自动清理时间(秒) |
FailedRetryInterval | 60(秒) | 失败消息重试轮询间隔 |
FailedRetryCount | 50 | 失败最大重试次数,超过后标记为永久失败 |
ConsumerThreadCount | 1 | 从消息队列消费的并发线程数 |
EnablePublishParallelSend | false | 是否用线程池并行执行发布 |
EnableSubscriberParallelExecute | false | 是否用内存队列并行执行订阅方法 |
SubscriberParallelExecuteThreadCount | Environment.ProcessorCount | 订阅并行执行的工作线程数 |
SubscriberParallelExecuteBufferFactor | 1 | 内存缓冲容量 = 线程数 × 该系数,提供背压保护 |
CollectorCleaningInterval | 300(秒) | 过期消息清理处理器运行间隔 |
FallbackWindowLookbackSeconds | 240(秒) | 重试处理器回看时间窗,容忍时钟偏移 |
SchedulerBatchSize | 1000 | 单个调度周期批量取出延迟/失败消息的最大数 |
UseStorageLock | false | 集群部署时是否用存储分布式锁保证单实例重试 |
配置方式统一在AddCap中完成,例如:
services.AddCap(x => { x.DefaultGroup = "my-default-group"; x.UseSqlServer("Your ConnectionString"); x.UseRabbitMQ("HostName"); x.FailedRetryCount = 10; });快速开始:InMemory 组合体验完整链路
介绍文档强调"对于新手非常友好",快速开始 提供了零外部依赖的启动路径:使用基于内存的事件存储和消息队列,一个控制台应用即可跑通发布-订阅全流程。
PM> Install-Package DotNetCore.CAP PM> Install-Package DotNetCore.CAP.InMemoryStorage PM> Install-Package Savorboard.CAP.InMemoryMessageQueuepublic void ConfigureServices(IServiceCollection services) { services.AddCap(x => { x.UseInMemoryStorage(); x.UseInMemoryMessageQueue(); }); }非 Web 的控制台程序需要手动引导启动与发布循环,参见 Sample.ConsoleApp/Program.cs:
var container = new ServiceCollection(); container.AddLogging(x => x.AddConsole()); container.AddCap(x => { x.UseInMemoryStorage(); x.UseInMemoryMessageQueue(); }).AddSubscribeFilter<Filter>(); var sp = container.BuildServiceProvider(); // 引导 CAP 后台处理器 sp.GetService<IBootstrapper>().BootstrapAsync(cts.Token); // 定时发布 await sp.GetService<ICapPublisher>().PublishAsync("sample.console.showtime", DateTime.Now, cancellationToken: cts.Token);该示例同时演示了订阅过滤器(SubscribeFilter的OnSubscribeExceptionAsync)的用法,用于在订阅异常时统一处理或重新抛出,对应 Filter 目录 的实现。
可靠性机制与监控运维
- 消息可靠性:CAP 内部会将消息持久化存储,配合重试(
IProcessor.NeedRetry等处理器,见 Processor 目录)与状态机流转,达到服务间数据最终一致性;异步消息隔离了故障传播,系统某部分的故障不会拖垮整个系统。 - 失败阈值回调:
FailedThresholdCallback允许在消息重试达到上限时介入处理(如告警)。 - 实时 Dashboard:安装
DotNetCore.CAP.Dashboard后默认通过/cap路径查看发布/接收消息及状态,可手动重试;分布式环境可结合 Consul(consul.md)或 Kubernetes(DotNetCore.CAP.Dashboard.K8s,kubernetes.md)进行节点发现。 - 可观测性:
DotNetCore.CAP.OpenTelemetry提供内置分布式追踪埋点。
CAP 整体架构:本地消息表与消息队列的协同,实现可靠投递与最终一致。
相关学习资料
介绍文档同时整理了系列学习资源(原链接为站外视频与博客,此处仅列出主题,可按名称检索):
- 视频教程:Bilibili、YouTube、腾讯视频上均有 CAP 系列入门教程;
- 文章系列:CAP 介绍及使用;CAP 7.0 / 6.0 / 5.0 / 3.0 / 2.6 / 2.5 / 2.4 / 2.3 各版本新特性解读;
- 社区里程碑:作为 .NET Core Community(NCC)早期千星项目,CAP 的成长记录见相关社区文章。
总结
CAP 把"事件总线"与"分布式事务"合二为一:对外提供约定大于配置、零侵入的发布订阅体验;对内通过本地消息表(Outbox 模式)+ 后台投递 + 自动重试,保证事件消息不丢失,实现微服务间的最终一致性。模块化的传输/存储/序列化抽象使其可以无缝对接 RabbitMQ、Kafka、Azure Service Bus 等消息队列与 SqlServer、MySql、PostgreSql、MongoDB 等数据库。对于正在构建微服务或 SOA 系统、希望在不引入复杂中间件的前提下获得可靠异步通信能力的团队,CAP 是一个开箱即用、可深度定制的务实选择。
- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
相关推荐
深入解析 CAP:基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案
深入解析 CAP:基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案 导读 CAP 是 .NET 社区(NCC)出品的分布式事务与事件总线
后端消息队列微服务CAP:基于本地消息表(Outbox 模式)的 .NET 分布式事务解决方案与事件总线
CAP:基于本地消息表(Outbox 模式)的 .NET 分布式事务解决方案与事件总线 CAP 是面向 .NET 平台的轻量级事件总线与分布式事务解决方案,它通
后端消息队列微服务消息路由CAP 事件总线与分布式事务解决方案:基于 Outbox 模式的最终一致性架构与实践
CAP 事件总线与分布式事务解决方案:基于 Outbox 模式的最终一致性架构与实践 CAP(Consistency And Partition)是 .NET
后端消息队列微服务
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考