.NET 9 实战:CAP 框架与 API 集成实现分布式事务最终一致性

📅 发布时间:2026/10/11 23:21:48
.NET 9 实战:CAP 框架与 API 集成实现分布式事务最终一致性
1. 项目背景CAP 与 API 集成的意义1.1 CAP 到底是什么如果你接触过微服务或分布式系统一定听过“分布式事务”这四个字。传统的做法是两阶段提交2PC但它在高并发、跨服务、跨数据库的场景下往往是灾难锁等待、协调者单点、性能损耗。而 CAP基于本地消息表的最终一致性方案走的是另一条路——把“分布式事务”拆成“本地事务 消息可靠投递”让数据最终一致而不是强一致。打个比方你在电商平台下单扣库存的系统和订单系统不在同一个数据库。用 CAP 的话订单系统在自己的数据库里先写入一条“用户创建了订单”的事件提交本地事务的同时把这条消息发到消息队列RabbitMQ / Kafka 等库存系统从队列里消费这条消息去做库存扣减。整个过程没有全局锁订单系统只需要保证“本地写入订单 写入消息”是原子的。只要这一步对了后续的消息投递由 CAP 负责重试和确认最终两边数据一定会对齐。所以 CAP 本质上是一个带持久化的发布订阅框架它处理了本地事务与消息发送的一致性、消息的可靠存储、重试、以及消费者幂等。.NET 生态里最主流的实现就是 DotNetCore.CAP目前社区活跃度很高.NET 9 下也能直接使用。1.2 为什么在 .NET 9 中要考虑“API 集成”网上很多教程只讲 CAP 的配置和发布订阅但实际项目里 CAP 几乎不可能独立存在。它要么被 API 层驱动——用户请求打到 Controller业务完成后发布事件要么在订阅方去调用别的服务的 API——跨系统通知、回调第三方平台、同步数据。这两种场景正好对应“API 集成”的含义。我写这篇文章就是想串起一条完整的链路从创建 .NET 9 项目到配置 CAP 连接消息队列和数据库到在 API 接口里发布事件再到订阅方通过 HTTP API 或消息体把结果写回业务系统。同时把那些“文档里不会写”的坑和排查思路一并放进来比如消息重复消费、401 鉴权失败、序列化不兼容、仪表盘怎么用等等。不管你是刚接触分布式事务还是已经在生产环境踩过坑这篇都值得花几分钟过一遍。2. 环境准备与基础设施选型2.1 .NET 9 开发环境搭建第一步自然是装 SDK。建议直接用当前版本的 .NET 9 SDK8.x 的项目也可以升级但既然标题写 .NET 9就直接以 9 为准。IDE 用 Visual Studio 2022 或者 Rider 都行命令行用dotnet --version确认版本dotnet --versionCAP 没有强依赖 .NET 9 专属的 API所以整个迁移成本很低。你只需要保证项目目标框架是net9.0然后正常 Install-Package 即可。不过有一点要提醒.NET 9 的空项目模板默认启用顶层语句top-level statements。这本身没问题但如果你习惯把 CAP 的配置写成一个扩展类注意builder的作用域。我通常把 CAP 注册单独抽成AddCapConfiguration保持Program.cs干净后续也方便做单元测试。2.2 消息队列RabbitMQ 还是 KafkaCAP 支持的消息队列非常多RabbitMQ、Kafka、Azure Service Bus、Amazon SQS、Redis Streams 等。做集成指南的话RabbitMQ 是首选因为部署简单、功能齐全、日志直观个人开发和中小团队足够用。Kafka 则是吞吐量优先适合超大规模事件流但运维复杂度明显高一个档次。我的建议是项目第一阶段先用 RabbitMQ把 CAP 的本地消息表跑顺再考虑是否迁移 Kafka。迁移时 CAP 的 API 基本不变只需要改UseRabbitMQ为UseKafka消费者代码完全不用动。这也是 CAP 的一个优势——消息队列对业务代码是抽象的。本地开发用 Docker 启动 RabbitMQ 非常方便docker run -d --name cap-rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERguest \ -e RABBITMQ_DEFAULT_PASSguest \ rabbitmq:3-management15672是管理界面端口5672是 AMQP 协议端口。默认账号密码都是guest生产环境一定要改。2.3 数据库SQL Server、PostgreSQL 还是 MySQLCAP 需要一张本地消息表来持久化事件数据所以必须有关系型数据库。不管你用哪种它都实时生成两张表cap.published已发布消息和cap.received已接收消息表名可按需修改。如果你所在团队已经在用 SQL Server直接用DotNetCore.CAP.SqlServer用 PostgreSQL 就配DotNetCore.CAP.PostgreSqlMySQL 用DotNetCore.CAP.MySql。以我实际经验PostgreSQL 做事件日志表最舒服占空间小、清理性能好配合 RabbitMQ 是常见组合。下面所有示例以 PostgreSQL 为准。docker run -d --name cap-postgres \ -e POSTGRES_PASSWORD123456 \ -e POSTGRES_DBcap_demo \ -p 5432:5432 \ postgres:162.4 项目结构设计既然要走“API 集成”这条路项目结构就不能太随意。我常用的是按职责分 ProjectDemo.Api包含 Controller、集成事件发布入口、CAP 消费者也可以放这里Demo.Application业务逻辑层接口实现、DTO 定义Demo.InfrastructureEF Core DbContext、仓储实现、CAP 配置Demo.Domain领域模型、事件定义如果你的项目规模不大三层结构就够了核心原则是事件模型定义放公共层生产者和消费者不要重复定义。CAP 的订阅默认根据消息类型反序列化如果两边类型不一致字段对不上就很容易出现“消息到了但是属性全是默认值”的诡异问题。3. 核心实现NuGet 包、注册配置与基础流程3.1 安装 CAP 相关 NuGet 包主包是DotNetCore.CAP另外依据消息队列和数据库选择安装额外包。操作上可以直接在Demo.Api项目里执行dotnet add package DotNetCore.CAP dotnet add package DotNetCore.CAP.RabbitMQ dotnet add package DotNetCore.CAP.PostgreSql dotnet add package DotNetCore.CAP.DashboardDashboard 包不是必要项但强烈建议装上。它给你提供一个 Web 页面可以直接查看消息状态、失败原因、重试次数后面排错时会发现它特别有用。注意包版本的一致性。尽量都选最新稳定版避免出现DotNetCore.CAP是 8.x、而RabbitMQ包是 7.x 的版本错位。用dotnet list package可以快速检查。3.2 Program.cs 中完成 CAP 服务注册在Program.cs里加注册代码。我一般放在 AddDbContext 之后因为它依赖数据库连接配置builder.Services.AddCap(x { // 使用 PostgreSQL 持久化消息状态 x.UsePostgreSql(Hostlocalhost;Databasecap_demo;Usernamepostgres;Password123456); // 使用 RabbitMQ 作为消息队列 x.UseRabbitMQ(options { options.HostName localhost; options.UserName guest; options.Password guest; options.Port 5672; options.VirtualHost /; }); // 全局默认分组名 x.DefaultGroupName cap.demo.api; // 失败重试次数 x.FailedRetryCount 3; // 成功消息保留时间 x.SucceedMessageExpiredAfter 24 * 3600; // 开启仪表盘 x.UseDashboard(); });这里每个参数都有讲究。DefaultGroupName是消费者分组名不同服务要用不同的名字否则多个服务会抢同一条消息。FailedRetryCount默认是 3但实际生产环境我习惯设成 5 甚至更高因为外部 API 偶发故障很常见重试次数太低会导致消息积压到失败队列还要人工去 Dashboard 里手动重发。SucceedMessageExpiredAfter控制成功消息在数据库里保留多久。默认 24 小时如果业务需要审计可以适当延长但要注意清理任务别让cap.published无限增长。3.3 初始数据库结构说明CAP 服务启动后会自动在数据库里创建消息表。如果你用的是 PostgreSQL执行成功以后可以看到两张表。这里简单解释它们的作用cap.published记录了每条已经发布的消息的 ID、消息体、队列名称、状态Succeeded / Failed、失败次数、最后重试时间cap.received记录消费者从队列里拿到的每条消息以及是否消费成功它们本质上就是原生消息队列的上游和下游落盘。一旦 RabbitMQ 中间件本身出现问题CAP 也能借助这两张表做恢复。这也是为什么“本地消息表 消息队列”要搭配数据库的原因消息队列保证高可用数据库保证可靠记录。如果你发现表没有生成先检查数据库连接串有没有权限建表。PostgreSQL 用户如果权限受限需要提前执行授权GRANT ALL PRIVILEGES ON SCHEMA public TO postgres;4. 实操从 API 接口发布事件到订阅方调用外部 API4.1 定义事件模型先定义一个集成事件Integration Event。这个类要放在公共的领域层或单独的契约类库中保证生产者和消费者引用的是同一份定义public record OrderCreatedEvent( Guid OrderId, string CustomerName, decimal TotalAmount, DateTime OccurredAt);用record类型很舒服因为 CAP 内部用 JSON 序列化record 的属性热重构不容易出错。字段枚举和时区也尽量在这个阶段约定好否则到了订阅端解析时时差或大小写问题会让你排查半天。4.2 API 接口里发布事件假设我们在Demo.Api里有一个下单接口。常规流程是写入订单表 - 发布OrderCreatedEvent- 返回成功。为了确保“本地数据库写入”和“消息发布”在同一个本地事务里完成需要显式开启事务并且用 CAP 提供的ICapPublisher[ApiController] [Route(api/orders)] public class OrdersController : ControllerBase { private readonly DemoDbContext _dbContext; private readonly ICapPublisher _capPublisher; public OrdersController(DemoDbContext dbContext, ICapPublisher capPublisher) { _dbContext dbContext; _capPublisher capPublisher; } [HttpPost] public async TaskIActionResult CreateOrder(CreateOrderRequest request) { // 开启事务 await using var transaction await _dbContext.Database.BeginTransactionAsync(); var order new Order { Id Guid.NewGuid(), CustomerName request.CustomerName, TotalAmount request.TotalAmount, CreatedAt DateTime.UtcNow }; _dbContext.Orders.Add(order); await _dbContext.SaveChangesAsync(); // 发布事件注意使用同一个事务 await _capPublisher.PublishAsync(order.created, new OrderCreatedEvent( order.Id, order.CustomerName, order.TotalAmount, order.CreatedAt)); await transaction.CommitAsync(); return Accepted(new { order.Id }); } }这里最关键的一行是_capPublisher.PublishAsync。它会把这条消息先写入cap.published表而不是立刻发给 RabbitMQ然后事务提交后CAP 后台进程再把消息真正发到交换机。这样即使 RabbitMQ 此刻宕机消息也不会丢因为它已经在本地库里落盘了。4.3 订阅方处理事件并调用外部 API订阅方可以是同一个 API 进程也可以是独立的服务。这里展示同进程内订阅方便演示。在Program.cs注册消费者时需要把订阅类注册为 DI 服务builder.Services.AddScopedOrderSubscriber();然后再写订阅类public class OrderSubscriber : ICapSubscribe { private readonly IHttpClientFactory _httpClientFactory; public OrderSubscriber(IHttpClientFactory httpClientFactory) { _httpClientFactory httpClientFactory; } [CapSubscribe(order.created)] public async Task HandleOrderCreated(OrderCreatedEvent event) { // 业务处理调用库存服务 API var client _httpClientFactory.CreateClient(StockService); var response await client.PostAsJsonAsync(/api/stock/deduct, new { event.OrderId, event.TotalAmount }); // 如果返回异常抛出异常让 CAP 触发重试 response.EnsureSuccessStatusCode(); } }消费者里如果调用了外部 API必须注意异常处理。CAP 默认会捕获异常并根据FailedRetryCount重试。这里一个常见误区是在订阅方法里用 try-catch 把异常吞掉CAP 会认为消息处理成功外部 API 没调到就假装没事了数据最终不对齐后面查问题非常痛苦。我的经验是正常情况下不要吞异常让它抛出去让 CAP 的重试机制去处理。4.4 事务与消息发布的原子性在上面的下单接口中可能有人会问“不是开发数据库事务了吗为什么消息发布成功提交数据库事务还没提交时消费者就已经拿到了消息”这个担心是有道理的。CAP 的处理方式是PublishAsync只是把消息写到本地消息表此时 RabbitMQ 并没有收到消息事务CommitAsync成功之后CAP 后台队列内部称为Publisher才会把消息表里状态为待发送的数据真正推送到 RabbitMQ。所以整个链路的时序是数据库事务提交订单写入orders表消息写入cap.published表CAP 后台线程轮询cap.published把消息发送到 RabbitMQRabbitMQ 把消息路由到消费者的队列消费者从cap.received记录状态并执行订阅方法如果第 1 步事务失败cap.published里的消息也不存在了不会出现“有事件但没业务数据”的情况。这就是典型的本地消息表模型。我遇到过不少新手在 EF Core 提交后忘写CapPublisher.PublishAsync结果业务数据有了、事件没发下游一直不更新。这是非常隐蔽的问题所以代码审查时最好盯着事务块内是否同时包含SaveChanges和PublishAsync。5. 常见问题与排查技巧实录5.1 调用外部 API 时返回 401 Unauthorized / incorrect api key最近看到很多社区讨论都提到“unexpected status 401 unauthorized: incorrect api key provided”这类错误虽然具体场景多半是调用大模型 API但在 CAP 集成外部业务 API 时同样常见。根本原因只有两类API Key 没传对或请求头没带上认证信息。排查步骤我建议按这个顺序来先抓实际发出的请求。在 HttpClient 里打日志确认 Authorization 头是否真的带上了而不是只在配置里写对了。确认 API Key 本身是否有效。很多平台创建 Key 后会有短暂延迟刚创建完立刻调用容易 401。确认网络出口 IP 是否在白名单。有些系统会校验来源 IP一旦环境切换本地到服务器就失败。确认键名。有的接口要求Authorization: Bearer key有的要求x-api-key: key写错一个位置也是 401。对应到 CAP 消费者里我建议把 API 凭据的读取抽出来放在配置中心或环境变量里并且配置变更要支持动态刷新。否则一旦密钥轮换所有消费者实例都要重启生产环境会很难看。另外要留意 CAP 重试机制和 401 的组合坑如果 401 是永久性错误比如 Key 确实错了CAP 重试 5 次也只是浪费资源消息还会一直停在失败队列。建议先检查cap.published对应消息的FailedCount如果全是 401先把配置改对再去 Dashboard 手动重发。5.2 消息重复消费的问题CAP 保证的是“至少一次”at least once不是“恰好一次”。也就是说在极端情况下——消费者处理完业务但还没提交确认时就宕机或者网络抖动导致 ACK 丢失——同一条消息可能被投递两次。解决方法只有一个核心思想消费者要做好幂等。最常见的手段是业务表加唯一索引。比如我们刚才的订单事件如果在订阅方要创建一张“库存扣减记录”那就给OrderId加唯一约束第二次消费时直接冲突捕获异常当成已处理。try { _dbContext.StockOperations.Add(stockOp); await _dbContext.SaveChangesAsync(); } catch (DbUpdateException ex) when (ex.InnerException is PostgresException pg pg.SqlState 23505) { // 唯一键冲突说明消息重复直接忽略 }切记不要在消费者里只做“判断存在与否”的查询然后跳过因为高并发时两个请求同时读到不存在还是可能重复写入。用数据库约束兜底才最稳。5.3 消息丢失或者无法触发消费有些时候消息在 Dashboard 里显示成功但消费者没执行。我遇到的情况基本是这三种发布者和订阅者的消息名不一致。CAP 的PublishAsync(order.created, ...)和[CapSubscribe(order.created)]必须字符串完全一致多一个空格都不行。订阅方法没被注册成 DI 服务。AddScopedOrderSubscriber()少了CAP 扫描不到订阅类自然不消费。多个实例用同一个DefaultGroupName消息被另一个实例消费了。特别是在调试环境本地开了多个服务RabbitMQ 里有好几个消费者绑定同一个队列消息会随机分配。排查这类问题Dashboard 是最好的入口。打开/cap页面默认路由查看Received区域如果消息状态是Succeed但业务没跑说明消费者拿到了消息但业务代码有问题如果Received里根本没有这条消息说明消息被其他实例消费了或路由配置不对。5.4 通过 Dashboard 监控运行状态CAP Dashboard 默认路由是/cap端口跟着 API 项目走。你可以看到几个关键指标Published / Received 的今日消息数失败消息列表、重试次数已发布的延迟消息失败详情包括完整的异常堆栈生产环境启用 Dashboard 时记得加鉴权最少也要用一个中间件限制内网访问。之前见过有人直接把 Dashboard 暴露公网消息体里的业务数据一览无余这个太危险了。6. 生产落地建议与进阶扩展6.1 重试策略与失败降级要单独设计CAP 内置的重试机制是线性重试立即重试失败后再次重试中间间隔很短。如果你的消费者是调用外部第三方 API线性重试很可能在外部服务故障恢复前就把次数用完。我实际的做法是在消费者里自己实现“延迟重试”收到消息后如果外部 API 返回 5xx 或超时把消息通过ICapPublisher发送到一个自定义的延迟队列比如order.retry然后订阅这个方法时手动Task.Delay做退避重试几次后如果还是失败就发到死信队列或写入失败日志表同时告警。CAP 本身没有内置“指数退避”但你可以基于它灵活扩展。一个简单的示例public async Task HandleOrderCreated(OrderCreatedEvent event) { var delay _context.GetRetryCount() * 10; // 伪代码实际从 CAP 的消息元数据中读取 await Task.Delay(TimeSpan.FromSeconds(delay)); // 重试外部 API }生产环境的可靠性不是靠某一个组件的重试次数堆出来的而是靠整体链路设计。CAP 只管消息不丢API 调用是否成功还是得业务自己兜底。6.2 消息体的向前兼容事件发布后消费者可能不会同步更新。比如今天发布了OrderCreatedEvent字段有 4 个下个月生产端加了第 5 个字段。老版本消费者反序列化时如果严格模式打开就会直接报错。CAP 默认的序列化策略是宽容的JSON 多出来的字段不会导致反序列化失败。但你要注意不要随意变更字段类型或删字段。我建议在事件模型层加版本号约定例如OrderCreatedEventV2新字段不加在老事件上而是用新事件名。这样才能保证消息消费者平滑升级。如果你使用的是 Confluent Kafka 之类的 schema registry这个问题会好一些但对 RabbitMQ CAP 的组合来说事件版本化主要靠团队纪律。6.3 密钥存储与安全审计这个话题容易被忽略。CAP 的连接字符串里包含数据库密码和 RabbitMQ 密码API 集成时会用到第三方 API Key这些都属于敏感信息。建议开发环境用appsettings.Development.json user secrets测试/生产环境一律用环境变量或配置中心不要提交到 Git日志里禁止打印完整报文。如果一定要打把 Authorization 头和信用卡号、手机号之类抹掉你可以写一个简单的 DelegatingHandler 对 HttpClient 出站请求做日志脱敏。这个处理在排查 401 问题时尤其重要——你既想看请求头又不想把密钥打到日志里。6.4 对 CAP 与 API 集成场景再扩展一点到这里基础的“API 接口发布 - CAP 投递 - 订阅方调用外部 API”链路已经完整了。实际项目中你还可以把 CAP 当异步任务队列用比如用户下单后API 立即返回库存扣减、发票生成、积分发放全部走事件异步处理支付回调进来后通过 CAP 发布“支付成功”事件多个下游服务并行消费定时任务里产生一批数据更新事件通过 CAP 分发到不同的 API 客户端这种异步化改造能够明显提升接口响应时间。我曾经把一个下单接口从 800ms 优化到 150ms就是因为在接口里只做了“写订单发事件”其余逻辑全部下沉到订阅方异步执行。当然异步化也意味着你要接受“用户请求结束但业务还没完全完成”的事实所以对外接口建议返回202 Accepted前端轮询查询订单状态而不是同步等待.最后再分享一个我踩过的坑在同一个方法里直接await _capPublisher.PublishAsync之后立刻调用另一个服务查询数据库期望事件消费者已经把数据写好了。结果自然是查不到因为 CAP 是异步投递的事件消费者执行有时间差。任何依赖事件结果的同步查询都是反模式正确的做法是查询时只依赖当前服务的数据需要汇聚结果时就查最终一致后的视图或者接受“稍后可见”。如果你要把这套东西真正放到生产我强烈建议先在测试环境全面验证一下故障场景停掉 RabbitMQ看消息是否积压在cap.published恢复后是否自动续发停掉消费者看消息是否积压在队列消费者恢复后是否堆积消费。这几条链路都验证通畅CAP API 集成这一套方案才算真正落地。