ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

CAP 框架全解析:基于 Outbox 模式的微服务事件总线与分布式事务解决方案

CAP 框架全解析:基于 Outbox 模式的微服务事件总线与分布式事务解决方案 后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载CAPdotnetcore/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内置实现消息队列TransportITransport、IConsumerClient、IConsumerClientFactoryTransport 目录Kafka、RabbitMQ、Azure Service Bus、AmazonSQS、NATS、Pulsar、Redis Streams、InMemoryMessageQueue存储StorageIDataStorage、IStorageInitializerPersistence 目录SqlServer、MySql、PostgreSql、MongoDB、InMemoryStorage序列化ISerializerSerialization 目录默认基于System.Text.JsonJsonSerializerOptions可自定义监控IMonitoringApiMonitoring 目录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 Dictionarystring, 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实现请求-响应式交互。核心配置参数一览CapOptionsCapOptions见 CAP.Options.cs集中管理处理管线的关键行为以下为构造函数中确认的默认值配置项默认值作用Versionv1消息版本标识≤20 字符用于多实例/多版本隔离DefaultGroupNamecap.queue.{程序集名小写}默认消费者组名SucceedMessageExpiredAfter8640024h成功消息自动清理时间秒FailedMessageExpiredAfter129600015 天失败消息自动清理时间秒FailedRetryInterval60秒失败消息重试轮询间隔FailedRetryCount50失败最大重试次数超过后标记为永久失败ConsumerThreadCount1从消息队列消费的并发线程数EnablePublishParallelSendfalse是否用线程池并行执行发布EnableSubscriberParallelExecutefalse是否用内存队列并行执行订阅方法SubscriberParallelExecuteThreadCountEnvironment.ProcessorCount订阅并行执行的工作线程数SubscriberParallelExecuteBufferFactor1内存缓冲容量 线程数 × 该系数提供背压保护CollectorCleaningInterval300秒过期消息清理处理器运行间隔FallbackWindowLookbackSeconds240秒重试处理器回看时间窗容忍时钟偏移SchedulerBatchSize1000单个调度周期批量取出延迟/失败消息的最大数UseStorageLockfalse集群部署时是否用存储分布式锁保证单实例重试配置方式统一在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.csvar container new ServiceCollection(); container.AddLogging(x x.AddConsole()); container.AddCap(x { x.UseInMemoryStorage(); x.UseInMemoryMessageQueue(); }).AddSubscribeFilterFilter(); var sp container.BuildServiceProvider(); // 引导 CAP 后台处理器 sp.GetServiceIBootstrapper().BootstrapAsync(cts.Token); // 定时发布 await sp.GetServiceICapPublisher().PublishAsync(sample.console.showtime, DateTime.Now, cancellationToken: cts.Token);该示例同时演示了订阅过滤器SubscribeFilter的OnSubscribeExceptionAsync的用法用于在订阅异常时统一处理或重新抛出对应 Filter 目录 的实现。可靠性机制与监控运维消息可靠性CAP 内部会将消息持久化存储配合重试IProcessor.NeedRetry等处理器见 Processor 目录与状态机流转达到服务间数据最终一致性异步消息隔离了故障传播系统某部分的故障不会拖垮整个系统。失败阈值回调FailedThresholdCallback允许在消息重试达到上限时介入处理如告警。实时 Dashboard安装DotNetCore.CAP.Dashboard后默认通过/cap路径查看发布/接收消息及状态可手动重试分布式环境可结合 Consulconsul.md或 KubernetesDotNetCore.CAP.Dashboard.K8skubernetes.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 CommunityNCC早期千星项目CAP 的成长记录见相关社区文章。总结CAP 把事件总线与分布式事务合二为一对外提供约定大于配置、零侵入的发布订阅体验对内通过本地消息表Outbox 模式 后台投递 自动重试保证事件消息不丢失实现微服务间的最终一致性。模块化的传输/存储/序列化抽象使其可以无缝对接 RabbitMQ、Kafka、Azure Service Bus 等消息队列与 SqlServer、MySql、PostgreSql、MongoDB 等数据库。对于正在构建微服务或 SOA 系统、希望在不引入复杂中间件的前提下获得可靠异步通信能力的团队CAP 是一个开箱即用、可深度定制的务实选择。赞分享后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐深入解析 CAP基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案深入解析 CAP基于 Outbox 模式的 .NET 微服务分布式事务与事件总线解决方案 导读 CAP 是 .NET 社区NCC出品的分布式事务与事件总线后端消息队列微服务CAP基于本地消息表Outbox 模式的 .NET 分布式事务解决方案与事件总线CAP基于本地消息表Outbox 模式的 .NET 分布式事务解决方案与事件总线 CAP 是面向 .NET 平台的轻量级事件总线与分布式事务解决方案它通后端消息队列微服务消息路由CAP 事件总线与分布式事务解决方案基于 Outbox 模式的最终一致性架构与实践CAP 事件总线与分布式事务解决方案基于 Outbox 模式的最终一致性架构与实践 CAPConsistency And Partition是 .NET后端消息队列微服务上一篇npx skills 交互式安装指南答对三个问题装好你的第一个 AI 技能下一篇Windows Cleaner 实战一招解决C盘空间不足的开源系统优化工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表