资讯动态

基于 Outbox 模式的分布式事务与事件总线:CAP 库完整集成指南

发布时间:2026/9/29 5:47:02 来源:尧图企业网站定制
后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载CAPDotNetCore.CAP是一个面向 .NET 生态的分布式事务与事件总线EventBus解决方案它采用与当前业务数据库集成的本地消息表Outbox 模式来保证分布式系统调用过程中事件消息在任何情况下都不会丢失。本指南以官方中文 README 为主体结合仓库源码与示例工程带你完整走通 CAP 的安装、配置、发布、订阅与 Dashboard 监控全流程并深入讲解其底层实现原理。核心思想用本地消息表解决最终一致性在构建 SOA 或微服务系统的过程中服务之间通常需要通过事件进行集成但单纯使用消息队列并不能保证数据的最终一致性——因为写业务数据库和发消息这两个动作无法原子完成。CAP 采用和当前数据库集成的本地消息表方案将待发送事件与业务数据写入同一条数据库事务中再由后台处理器异步把消息投递到消息队列从而解决分布式系统互相调用各个环节可能出现的异常。CAP 实现了 eShop 电子书中描述的发件箱模式Outbox Pattern。同时CAP 也可以当作一个纯粹的 EventBus 使用它提供了一种更加简单的事件发布与订阅方式在订阅与发布过程中你不需要继承或实现任何接口订阅方法位于 Controller 中时。架构预览从上图可以看到 CAP 在微服务架构中的位置每个微服务内部的业务数据库都内置了一张 CAP 消息表业务数据与事件消息在同一事务中落库服务启动后CAP 的后台处理器会把消息表里的待发送消息投递到消息队列RabbitMQ、Kafka 等同时订阅端从消息队列拉取消息执行订阅方法。这一机制正是分布式事务最终一致性的落地基础。安装NuGet在你的项目中运行以下命令安装核心包PM Install-Package DotNetCore.CAPCAP 支持主流的消息队列作为传输器Transport按需选择安装PM Install-Package DotNetCore.CAP.Kafka PM Install-Package DotNetCore.CAP.RabbitMQ PM Install-Package DotNetCore.CAP.AzureServiceBus PM Install-Package DotNetCore.CAP.AmazonSQS PM Install-Package DotNetCore.CAP.NATS PM Install-Package DotNetCore.CAP.RedisStreams PM Install-Package DotNetCore.CAP.PulsarCAP 提供主流数据库作为存储Storage事件日志表会被集成到你选择的数据库中// 按需选择安装你正在使用的数据库 PM Install-Package DotNetCore.CAP.SqlServer PM Install-Package DotNetCore.CAP.MySql PM Install-Package DotNetCore.CAP.PostgreSql PM Install-Package DotNetCore.CAP.MongoDB注意MongoDB 存储仅支持 MongoDB 4.0 集群。配置AddCap 入门在Startup.cs或Program.cs中通过AddCap配置 CAP。仓库中的示例工程 Sample.RabbitMQ.SqlServer 给出了完整的配置形态public void ConfigureServices(IServiceCollection services) { ...... services.AddDbContextAppDbContext(); services.AddCap(x { // 如果你使用 EF 进行数据操作添加如下配置可选项无需再配置 x.UseSqlServer 了 x.UseEntityFrameworkAppDbContext(); // 如果你使用 ADO.NET根据数据库选择进行配置 x.UseSqlServer(数据库连接字符串); x.UseMySql(数据库连接字符串); x.UsePostgreSql(数据库连接字符串); // 如果你使用 MongoDB可以添加如下配置 x.UseMongoDB(ConnectionStrings); // 注意仅支持 MongoDB 4.0 集群 // CAP 支持 RabbitMQ、Kafka、AzureServiceBus、AmazonSQS 等作为 MQ根据使用选择配置 x.UseRabbitMQ(ConnectionStrings); x.UseKafka(ConnectionStrings); x.UseAzureServiceBus(ConnectionStrings); x.UseAmazonSQS(); }); }存储与传输的选择原则存储决定 CAP 消息表Published/Received 表落在哪里必须与你的业务数据库一致这样事件才能与业务数据同事务写入传输决定事件最终投递到哪个消息队列。二者可自由组合例如UseEntityFramework UseRabbitMQ、UseMySql UseKafka。核心配置项CapOptionsCapOptions 源码 给出了全部可调参数及其默认值下面是高频使用的几个配置项默认值说明DefaultGroupNamecap.queue.{入口程序集名小写}默认消费者组名在 Kafka 中对应消费者组group.id在 RabbitMQ 中对应队列名GroupNamePrefix/TopicNamePrefix空为所有消费者组名 / Topic 名增加可选前缀Versionv1消息版本标识最长 20 字符用于隔离不同实例或部署之间的数据SucceedMessageExpiredAfter8640024 小时成功处理后的消息自动删除的时间间隔秒FailedMessageExpiredAfter129600015 天超过重试阈值后的失败消息自动清理时间间隔秒FailedRetryInterval60重试处理器轮询并重试失败消息的时间间隔秒FailedRetryCount50失败消息发布与订阅的最大重试次数超过后标记为永久失败不再重试FailedThresholdCallback空消息重试达到上限仍失败时的回调可在此做告警或人工介入ConsumerThreadCount1从传输器消费消息的并发消费者线程数EnableSubscriberParallelExecutefalse是否启用订阅方法的并行执行内存队列缓冲 多线程SubscriberParallelExecuteThreadCount逻辑处理器数并行执行时的 Worker 线程数EnablePublishParallelSendfalse是否用线程池并行执行发布任务CollectorCleaningInterval300清理过期消息处理器运行间隔秒FallbackWindowLookbackSeconds240重试处理器回溯获取 Scheduled/Failed 状态消息的时间窗秒兼容时钟轻微偏差SchedulerBatchSize1000单次调度周期抓取的延迟/失败消息批量大小UseStorageLockfalse是否使用分布式存储锁执行重试集群部署时避免多个实例重复重试JsonSerializerOptions默认消息内容序列化/反序列化选项可定制命名策略、转换器等例如设置默认消费者组与失败重试回调services.AddCap(x { x.DefaultGroupName my-default-group; x.FailedRetryCount 20; x.FailedThresholdCallback failed { // 记录日志、告警或人工排查 var logger failed.ServiceProvider.GetRequiredServiceILoggerProgram(); logger.LogError($A message of type {failed.MessageType} failed after retrying {x.FailedRetryCount} times, requiring manual troubleshooting. Message name: {failed.Message.GetName()}); }; });链式扩展CapBuilderAddCap返回的 CapBuilder 还提供两个实用扩展方法AddSubscribeFilterT()注册订阅过滤器需实现ISubscribeFilter用于日志、异常处理、事务管理等横切关注点按注册顺序执行作用域Scoped生命周期AddSubscriberAssembly(params Assembly[] | params Type[])指定程序集扫描订阅方法默认情况下 CAP 会扫描所有已加载程序集该方法可限定扫描范围以提升启动性能。发布消息在 Controller 中注入ICapPublisher然后调用Publish系列方法发送消息。仓库示例 ValuesController.cs 覆盖了全部发布形态public class PublishController : Controller { private readonly ICapPublisher _capBus; public PublishController(ICapPublisher capPublisher) { _capBus capPublisher; } // 不使用事务 [Route(~/without/transaction)] public IActionResult WithoutTransaction() { _capBus.Publish(xxx.services.show.time, DateTime.Now); // 发布延迟消息版本 7.0 支持发送延迟消息 _capBus.PublishDelayAsync(TimeSpan.FromSeconds(delaySeconds), xxx.services.show.time, DateTime.Now); return Ok(); } // Ado.Net 中使用事务自动提交 [Route(~/adonet/transaction)] public IActionResult AdonetWithTransaction() { using (var connection new MySqlConnection(ConnectionString)) { using (var transaction connection.BeginTransaction(_capBus, autoCommit: true)) { //业务代码 _capBus.Publish(xxx.services.show.time, DateTime.Now); } } return Ok(); } // EntityFramework 中使用事务自动提交 [Route(~/ef/transaction)] public IActionResult EntityFrameworkWithTransaction([FromServices]AppDbContext dbContext) { using (var trans dbContext.Database.BeginTransaction(_capBus, autoCommit: true)) { //业务代码 _capBus.Publish(xxx.services.show.time, DateTime.Now); } return Ok(); } }关于事务的关键理解BeginTransaction(_capBus, autoCommit: true)会把 CAP 消息写入和业务数据写入纳入同一条数据库事务——事务提交时业务数据与事件消息一起落库事务回滚则一起撤销从根本上避免数据已写、消息未发或反之的不一致问题。autoCommit: false时则需要手动Commit()参见示例中AdonetWithTransaction的BeginTransactionAsyncCommitAsync用法。ICapPublisher API 全览ICapPublisher 接口 提供以下发布方法均有同步与异步版本PublishT(name, contentObj, callbackName)/PublishT(name, contentObj, headers)发布消息name为 Topic 名或交换机路由键contentObj会被序列化可为 null可选callbackName回调订阅名或附加自定义headersPublishDelayAsyncT(delayTime, name, contentObj, ...)发布延迟消息延迟由 CAP 内部调度实现不依赖消息队列的延迟队列特性Transaction属性暴露当前 CAP 事务上下文供事务性发布使用。订阅消息Action MethodController 内订阅在 Action 上添加CapSubscribeAttribute来订阅相关消息public class PublishController : Controller { [CapSubscribe(xxx.services.show.time)] public void CheckReceivedMessage(DateTime datetime) { Console.WriteLine(datetime); } }Service Method业务服务中订阅如果你的订阅方法不在 Controller 中订阅类需要继承ICapSubscribe接口namespace xxx.Service { public interface ISubscriberService { void CheckReceivedMessage(DateTime datetime); } public class SubscriberService: ISubscriberService, ICapSubscribe { [CapSubscribe(xxx.services.show.time)] public void CheckReceivedMessage(DateTime datetime) { } } }然后在ConfigureServices()中注入你的ISubscriberService类public void ConfigureServices(IServiceCollection services) { services.AddTransientISubscriberService,SubscriberService(); services.AddCap(x{}); }异步订阅订阅方法可以返回Task并接收CancellationToken参数实现异步订阅public class AsyncSubscriber : ICapSubscribe { [CapSubscribe(name)] public async Task ProcessAsync(Message message, CancellationToken cancellationToken) { await SomeOperationAsync(message, cancellationToken); } }使用多部分订阅名Partial Topic如果想在类级别对订阅的 Topic 进行分组可以将方法上的订阅设置为部分订阅isPartial: true消息队列上的最终订阅名将是类上的 topic 方法上的 topic的拼合。下面示例中当收到customers.create消息时将调用Create(..)函数[CapSubscribe(customers)] public class CustomersSubscriberService : ICapSubscribe { [CapSubscribe(create, isPartial: true)] public void Create(Customer customer) { } }订阅者组Group订阅者组的概念类似于 Kafka 中的消费者组与消息队列中的广播模式相同用来处理不同微服务实例之间同时消费相同消息的场景。从 TopicAttribute 源码 可以看到Group和GroupConcurrent是特性上的两个关键属性默认行为CAP 启动时会创建一个默认消费者组DefaultGroupName。如果多个相同消费者组的消费者订阅同一个 Topic 消息只会有一个消费者被执行竞争消费如果消费者位于不同的消费者组则所有消费者都会被执行广播/扇出。指定分组同一实例中可以通过Group参数指定不同的消费者组[CapSubscribe(xxx.services.show.time, Group group1 )] public void ShowTime1(DateTime datetime) { } [CapSubscribe(xxx.services.show.time, Group group2)] public void ShowTime2(DateTime datetime) { }ShowTime1和ShowTime2将被同时调用。GroupConcurrent限制该订阅并发消费的消息数量如果设置了该值但未指定GroupCAP 会自动使用Name创建 Group源码注释。默认组名配置services.AddCap(x { x.DefaultGroup default-group-name; });Dashboard 仪表盘CAP 提供了仪表盘Dashboard功能可以很方便地查看发出和接收到的消息及其状态并实时查看消息收发情况。安装命令PM Install-Package DotNetCore.CAP.Dashboard仪表盘默认访问地址是http://localhost:xxx/cap你可以在d.MatchPath配置项中把cap路径后缀修改为其他名字对应UseDashboard的PathMatch选项。分布式节点查看在分布式环境中仪表盘内置集成了 Consul 作为节点的注册发现同时实现了网关代理功能你可以像访问本地资源一样方便地查看本节点或其他节点的数据。Consul 发现配置详见 Consul 配置文档核心配置形如services.AddCap(x { x.UseMySql(Configuration.GetValuestring(ConnectionString)); x.UseRabbitMQ(localhost); x.UseDashboard(); x.UseConsulDiscovery(_ { _.DiscoveryServerHostName localhost; _.DiscoveryServerPort 8500; _.CurrentNodeHostName Configuration.GetValuestring(ASPNETCORE_HOSTNAME); _.CurrentNodePort Configuration.GetValueint(ASPNETCORE_PORT); _.NodeId Configuration.GetValuestring(NodeId); _.NodeName Configuration.GetValuestring(NodeName); }); });Kubernetes 部署如果你的服务部署在 Kubernetes 中请使用为 Kubernetes 专门提供的发现包PM Install-Package DotNetCore.CAP.Dashboard.K8s配置方式详见 Kubernetes 配置文档。深入源码CAP 的底层设计消息发布与事务的默认实现核心包 ICapPublisher.Default.cs 是ICapPublisher的默认实现负责把发布请求包装为 CAP 消息并写入存储IDataStorage 接口 定义了消息的持久化契约各数据库包SqlServer / MySql / PostgreSql / MongoDB / InMemoryStorage分别实现该接口因此事件表能无缝集成进你的业务数据库。后台处理器保证消息必达Processor 目录 下的多个处理器IProcessingServer.Cap.cs、IProcessor.NeedRetry.cs、IProcessor.Delayed.cs、IProcessor.Collector.cs等构成了消息可靠性投递的引擎消息写入本地表后处理器异步将其发送到 MQ发送失败的进入重试队列按FailedRetryInterval轮询重试最多FailedRetryCount次延迟消息由ScheduledMediumMessageQueue调度到时间点后投递成功/过期消息由清理器按SucceedMessageExpiredAfter/FailedMessageExpiredAfter定期回收。这也是任何情况下事件消息都不会丢失的实现保障。订阅扫描与调用订阅方法的发现由 IConsumerServiceSelector 完成MethodMatcherCache负责缓存方法匹配结果真正执行订阅方法的是 ISubscribeInvoker。CapSubscribeAttributeCAP.Attribute.cs继承自 TopicAttribute后者支持AttributeTargets.Method | AttributeTargets.Class这正是类级与方法级组合出多部分订阅名Partial Topic的机制来源。快速体验仓库自带示例仓库samples/目录提供了多种组合的可运行示例是最佳的实践参考Sample.RabbitMQ.SqlServerRabbitMQ EF/SqlServerController 中完整展示了无事务发布、ADO.NET 事务发布、EF 事务发布、延迟发布、延迟事务发布以及CapSubscribe订阅见 ValuesController.cs并附带数据库迁移Migrations/与 TypedConsumer 扩展用法Sample.RabbitMQ.MySql、Sample.RabbitMQ.MongoDB含 docker-compose其他存储组合的对照实现Sample.Kafka.PostgreSqlKafka PostgreSql 组合Sample.ConsoleApp非 Web 场景下使用内存存储 内存队列 订阅过滤器SubscribeFilter的轻量用法Sample.Dashboard.Jwt / Sample.Dashboard.AuthDashboard 自定义认证与 JWT 鉴权示例。贡献与 License欢迎通过参与讨论、报告 issue 或提交 Pull Request 为 CAP 贡献代码。项目采用 MIT License 开源见 LICENSE.txt你可以自由地在商业与个人项目中使用。小结CAP 用一条与业务数据库同事务的本地消息表把业务提交与事件投递原子化再借助后台处理器与重试机制保证消息最终必达同时它以零接口侵入的方式提供发布/订阅、延迟消息、消费者组、并行消费、Dashboard 与 OpenTelemetry 等能力。上手路径十分清晰安装传输器与存储包 →AddCap配置 → 注入ICapPublisher发布 → 打上[CapSubscribe]订阅 → 打开 Dashboard 观察消息流转。对照本仓库的源码与示例工程你可以在自己的微服务项目中快速落地这套最终一致性方案。赞分享后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐ESP-IDF Wi-Fi 场景低功耗模式DTIM 省电原理与 Modem-sleep / Auto Light-sleep / Deep-sleep 配置实战ESP IDF Wi Fi 场景低功耗模式DTIM 省电原理与 Modem sleep / Auto Light sleep / Deep sleep 配置实后端消息队列微服务CAP 分布式事务与事件总线实战指南基于 Outbox 模式的消息可靠投递与订阅.NETCAP 分布式事务与事件总线实战指南基于 Outbox 模式的消息可靠投递与订阅.NET 导读 CAP 是一个面向 .NET 的轻量级分布式事务与事件总线后端消息队列微服务消息路由CAP 分布式事务与事件总线实战指南基于 Outbox 模式的最终一致性消息解决方案CAP 分布式事务与事件总线实战指南基于 Outbox 模式的最终一致性消息解决方案 本文基于 .NET 开源库 CAPDotNetCore.CAP讲解如后端消息队列微服务消息路由上一篇告别Cron表达式噩梦vue-cron可视化构建工具全指南下一篇【亲测免费】 smzdm_script 项目安装和配置指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑