
后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载导读CAP 是 .NET 生态中一个同时扮演EventBus事件总线与分布式事务解决方案双重角色的框架专为微服务与 SOA 架构设计帮助你构建可扩展、可靠且易于演进的微服务系统。本文将以官方介绍文档为主线结合仓库源码系统讲解 CAP 的核心定位、EventBus 设计理念、模块化架构、默认配置参数以及一套可直接运行的快速开始示例让你既能理解为什么用 CAP也能马上上手怎么用 CAP。CAP 是什么EventBus 与分布式事务的统一体CAP 是一个 EventBus同时也是一个在微服务或者 SOA 系统中解决分布式事务问题的框架。它有助于创建可扩展、可靠并且易于更改的微服务系统。在微软的eShop微服务示例项目中CAP 被推荐作为生产环境可用的 EventBus。这也意味着 CAP 的定位并非实验性玩具而是可以承载真实业务流量的消息基础设施。它的核心价值来自两点作为 EventBus让不同服务之间以发布/订阅的方式解耦通信无需彼此感知对方的存在作为分布式事务方案基于Outbox发件箱模式 消息队列以最终一致性取代跨服务强一致的分布式事务避免引入分布式事务协调器的复杂度。从仓库源码看这两点分别由两组核心抽象支撑发布侧由 ICapPublisher 负责将消息送入事件总线订阅侧由 ICapSubscribe 与 CapSubscribeAttribute 负责自动发现并执行订阅方法。什么是 EventBus组件间互相不知道也能通信!!! question 什么是 EventBus事件总线是一种机制它允许不同的组件彼此通信而不彼此了解。组件可以将事件发送到 EventBus而无需知道是谁来接听或有多少其他人来接听。组件也可以侦听 EventBus 上的事件而无需知道谁发送了事件。这样组件可以相互通信而无需相互依赖。同样很容易替换一个组件——只要新组件了解正在发送和接收的事件其他组件就永远不会知道。这段定义点出了 EventBus 的三大收益解耦发送方与接收方之间不存在任何直接引用或调用关系弹性替换只要新组件理解相同的事件格式替换旧组件对其他组件完全透明故障隔离系统某一部分的故障不会直接蔓延至其他部分。CAP 在此之上更进一步它把发送事件从业务代码里直接调用消息队列客户端中解放出来。业务方法只需要发布一个对象POCO序列化、投递、确认、重试全部由框架接管。CAP 的设计特色无接口约束、约定大于配置、轻量灵活相对于其他的 Service Bus 或者 Event BusCAP 拥有自己的特色它不要求使用者发送消息或者处理消息的时候实现或者继承任何接口拥有非常高的灵活性。这一点在源码中得到充分印证发布消息时你只需要注入 ICapPublisher这是框架提供的基础服务业务侧无需实现任何发布接口调用Publish/PublishAsync即可消费消息时ICapSubscribe 只是一个空标记接口——它不定义任何成员唯一作用是让框架在启动时自动扫描并发现带[CapSubscribe]特性的订阅方法订阅方法本身是普通方法只要加上 CapSubscribeAttribute 特性参数可以是一个普通 DTO 对象甚至可以直接接收DateTime等基础类型。CAP 团队约定大于配置的理念使框架对于新手非常友好并且拥有轻量级不需要继承基类、不需要实现接口、不需要定义消息契约类层次写一个普通方法即可成为订阅者。模块化架构消息队列、存储、序列化皆可替换CAP 采用模块化设计具有高度的可扩展性。你有许多选项可以选择包括消息队列、存储、序列化方式等系统的许多元素内容可以替换为自定义实现。这种可扩展性在代码结构上体现得非常清晰见仓库 src 目录核心项目DotNetCore.CAP只定义抽象与默认流程具体实现全部以独立程序集的方式按需引入消息队列TransportRabbitMQ、Kafka、Azure Service Bus、Amazon SQS、NATS、Pulsar、Redis Streams以及面向本地开发/演示的 InMemory Queue存储StorageSQL Server、MySQL、PostgreSQL、MongoDB以及 InMemoryStorage序列化默认使用System.Text.Json可通过CapOptions.JsonSerializerOptions定制见 CAP.Options.cs也支持通过ISerializer替换为自定义实现监控与运维内置 Dashboard 可视化面板、OpenTelemetry 可观测性集成、Consul / K8s 集群节点发现。这种核心 扩展包的形态正是通过 ICapOptionsExtension 这一扩展点实现的第三方存储或消息队列只需实现AddServices(IServiceCollection)注册自己的IDataStorage、ITransport等实现即可无缝接入。而 CapBuilder 则进一步提供了AddSubscribeFilterT()订阅过滤器与AddSubscriberAssembly(...)指定程序集扫描等细粒度能力。架构与工作原理Outbox 模式如何保证最终一致性下图展示了 CAP 在微服务间实现分布式事务最终一致性的完整链路从图中可以看出核心流程分为四个环节本地事务写入微服务 A 执行业务 SQL并在同一个本地数据库事务中同时写入事件 SQL即 Outbox 表可靠投递CAP 组件将事件发送到消息队列RabbitMQ / Kafka 等消息中转消息队列将事件转发给微服务 B 对应的 CAP 组件事件消费微服务 B 的 CAP 组件在本地事务中写入事件 SQL 并执行对应业务逻辑。支撑这一流程的关键抽象包括ICapTransaction把消息发布与数据库事务绑定在一起提供Commit/Rollback保证数据库变更与消息发送同生共死——事务提交成功才真正投递消息事务回滚则丢弃缓冲的消息IDataStorage统一的消息存储接口包含StoreMessageAsync存储待发送消息、ChangePublishStateAsync更新消息状态、GetPublishedMessagesOfNeedRetry取出需要重试的消息等方法是重试与最终一致性的数据基础IStorageInitializer负责初始化存储结构GetPublishedTableName/GetReceivedTableName/GetLockTableName即框架会在数据库中自动创建发布消息表、接收消息表与分布式锁表消息自带标准头部见 Headers.cs包括cap-msg-id消息唯一 ID、cap-msg-name主题名、cap-msg-group消费组、cap-senttime发送时间、cap-delaytime延迟发送时间等用于路由、幂等、追踪与延迟调度。可靠性是 CAP 区别于直接集成消息队列的关键消息在发出前已持久化配合失败重试、延迟调度、过期清理等后台处理器见 Processor 目录下的IProcessor.NeedRetry、IProcessor.Delayed、IProcessor.Collector等实现系统某一部分的故障不会导致整个系统崩溃各服务数据最终达成一致。核心配置参数速查CapOptions 默认值CAP 的全局行为通过 CapOptions 配置下表整理了源码构造函数中给出的默认值与含义供你在调优时参考配置项默认值含义DefaultGroupNamecap.queue.{入口程序集名小写}默认消费组名Kafka 中对应消费者组RabbitMQ 中对应队列名GroupNamePrefix/TopicNamePrefixnull可选的消费组名 / 主题名前缀Versionv1消息版本号用于多实例/多版本共存时的数据隔离最长 20 字符SucceedMessageExpiredAfter8640024 小时处理成功消息的自动清理时间秒FailedMessageExpiredAfter129600015 天失败消息的自动清理时间秒FailedRetryInterval60重试处理器轮询失败消息的时间间隔秒FailedRetryCount50失败消息最大重试次数超过后标记为永久失败FailedThresholdCallback—达到最大重试次数时的回调可接收 FailedInfo 做告警ConsumerThreadCount1从消息队列拉取消息的消费线程数EnableSubscriberParallelExecutefalse是否启用订阅方法内存队列并行执行SubscriberParallelExecuteThreadCount逻辑处理器数量并行执行订阅方法的工作线程数SubscriberParallelExecuteBufferFactor1并行缓冲容量倍数缓冲容量 线程数 × 因子EnablePublishParallelSendfalse是否启用线程池并行发布消息CollectorCleaningInterval3005 分钟清理过期消息的后台处理器运行间隔秒FallbackWindowLookbackSeconds2404 分钟重试处理器回看调度/失败消息的时间窗口秒用于容忍时钟偏差SchedulerBatchSize1000单个调度周期批量取出延迟/失败消息的数量上限UseStorageLockfalse是否启用分布式存储锁集群部署时避免多实例重复重试JsonSerializerOptionsnew()消息内容 JSON 序列化选项快速开始五分钟跑通发布与订阅以下示例直接取自官方快速开始文档docs/content/user-guide/zh/getting-started/quick-start.md并补充了源码层面的说明。1. 安装 NuGet 包使用基于内存的事件存储和消息队列可零外部依赖快速启动PM Install-Package DotNetCore.CAP PM Install-Package DotNetCore.CAP.InMemoryStorage PM Install-Package Savorboard.CAP.InMemoryMessageQueue生产环境请将存储与消息队列替换为对应的扩展包如DotNetCore.CAP.RabbitMQ、DotNetCore.CAP.Kafka搭配DotNetCore.CAP.SqlServer/DotNetCore.CAP.MySql/DotNetCore.CAP.PostgreSql/DotNetCore.CAP.MongoDB详见仓库 src 目录。2. 在 ASP.NET Core 中注册 CAP在Startup.cs或.NET 6的Program.cs中添加配置public void ConfigureServices(IServiceCollection services) { services.AddCap(x { x.UseInMemoryStorage(); x.UseInMemoryMessageQueue(); }); }AddCap返回的 CapBuilder 还支持链式调用AddSubscribeFilterT()、AddSubscriberAssembly(...)等高级配置。3. 发送消息public class PublishController : Controller { [Route(~/send)] public IActionResult SendMessage([FromServices]ICapPublisher capBus) { capBus.Publish(test.show.time, DateTime.Now); return Ok(); } }这里调用的正是 ICapPublisher 中的PublishT(string name, T? contentObj, ...)重载第一个参数是主题名Topic / Exchange 路由键第二个参数是会被序列化的消息体对象。发送延迟消息public class PublishController : Controller { [Route(~/send/delay)] public IActionResult SendDelayMessage([FromServices]ICapPublisher capBus) { capBus.PublishDelay(TimeSpan.FromSeconds(100), test.show.time, DateTime.Now); return Ok(); } }延迟消息会通过cap-delaytime头部记录预定发送时间见 Headers.cs由延迟调度处理器在到达时间后投递最终进入消息队列。发送包含头信息的消息var header new Dictionarystring, string() { [my.header.first] first, [my.header.second] second }; capBus.Publish(test.show.time, DateTime.Now, header);自定义头部用于携带路由或业务元数据与 CAP 自带的标准头部cap-msg-id等相互独立消费端可通过[FromCap] CapHeader读取。4. 处理消息public class ConsumerController : Controller { [NonAction] [CapSubscribe(test.show.time)] public void ReceiveMessage(DateTime time) { Console.WriteLine(message time is: time); } }[CapSubscribe(test.show.time)]声明该方法订阅指定主题。CAP 在启动时通过 IConsumerServiceSelector 扫描程序集中带有该特性的方法并注册为消费者无需实现任何接口——这正是前文无接口约束特色的落地体现。处理包含头信息的消息[CapSubscribe(test.show.time)] public void ReceiveMessage(DateTime time, [FromCap]CapHeader header) { Console.WriteLine(message time is: time); Console.WriteLine(message first header : header[my.header.first]); Console.WriteLine(message second header : header[my.header.second]); }[FromCap]标记参数由框架注入见 CAP.Attribute.cs 中的FromCapAttributeCapHeader本质是ReadOnlyDictionarystring, string?可直接按下标读取自定义头部还可调用AddResponseHeader、RewriteCallback等方法处理请求-响应回调场景。5. 更进一步事务内发布消息在业务代码中结合数据库事务发布消息是 CAP 最终一致性的标准用法。注入ICapTransaction并将事务对象关联到ICapPublisher.Transaction即可保证业务数据与事件在同一个本地事务中提交详见 ICapTransaction.cs 与各存储包的ICapTransaction.*实现。摘要为什么 CAP 比直接集成消息队列更可靠相对于直接集成消息队列异步消息传递最强大的优势之一是可靠性系统的一个部分中的故障不会传播也不会导致整个系统崩溃。在 CAP 内部会将消息进行存储以保证消息的可靠性并配合重试等策略以达到各个服务之间的数据最终一致性。消息在发布前先持久化到 Outbox 表投递失败可自动重试默认最多 50 次、间隔 60 秒成功消息在 24 小时后自动清理失败消息在 15 天后自动清理存储不会无限膨胀集群部署时可开启UseStorageLock配合锁表避免多实例重复消费与重复重试配合内置 Dashboard见 src/DotNetCore.CAP.Dashboard可以可视化地查看发布/订阅消息的状态与重试情况。延伸阅读快速开始与完整示例docs/content/user-guide/zh/getting-started/quick-start.md详细配置说明docs/content/user-guide/zh/cap/configuration.md消息发送/订阅进阶docs/content/user-guide/zh/cap/messaging.md事务与本地消息表docs/content/user-guide/zh/cap/transactions.md消息队列接入RabbitMQ / Kafka / NATS 等docs/content/user-guide/zh/transport/存储接入SQL Server / MySQL / PostgreSQL / MongoDB 等docs/content/user-guide/zh/storage/源码级实现src/DotNetCore.CAP、src/DotNetCore.CAP.Dashboard赞分享后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载相关推荐深入解析 CAP基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案深入解析 CAP基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案 导读 CAP 是 .NET 社区NCC出品的分布式事务与事件总线后端消息队列微服务CAP基于本地消息表Outbox 模式的 .NET 分布式事务解决方案与事件总线CAP基于本地消息表Outbox 模式的 .NET 分布式事务解决方案与事件总线 CAP 是面向 .NET 平台的轻量级事件总线与分布式事务解决方案它通后端消息队列微服务消息路由CAP 框架全解析基于 Outbox 模式的微服务事件总线与分布式事务解决方案CAP 框架全解析基于 Outbox 模式的微服务事件总线与分布式事务解决方案 CAPdotnetcore/CAP是一个开箱即用的 .NET 事件总线E后端消息队列微服务上一篇如何快速上手OpenDesign DataStat5分钟搭建你的开源社区数据面板下一篇MAS 激活脚本完整教程4 条路线 10 分钟搞定 Windows 与 Office 永久激活创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考