
做后端开发这几年我越来越觉得系统越拆一致性越难。订单接口刚给前端返了 200库存服务那边却超时了钱扣了、单也下了商品却因为没扣到库存被别人买走。这种问题大多数人的第一反应是把接口改成同步调用、加超时重试但试过的人都知道治标不治本。真正稳妥的做法是引入 CAP 这样的分布式事务框架——把“业务操作”和“要发给别人的事件”放在同一个本地事务里事件先落到消息表再由后台可靠投递到消息队列下游订阅方按需消费。今天这篇就基于 .NET 9 从零搭一套 CAP 服务并把它接到真实的 API 链路里讲清楚每一步为什么这么写。1. 为什么是 CAP几种分布式事务方案里它最贴近日常 API 集成的真实需求1.1 核心原理把“说出去的话”先写进自己的本子CAP 的核心思路其实很简单当你需要通知另一个服务去做某件事时不要把“通知”这个动作当成随手的网络调用而是先把这件事记录在业务数据库的一张表里。这个模式和“本地消息表 / Outbox”几乎一样。用一句话概括业务写库和事件写库发生在同一个数据库事务里事务提交后CAP 的后台进程负责把事件安全地投递到消息队列。这样你的 API 接口不用关心下游服务是否在线、调用是否超时只要事务提交成功事件就“必然会被送出去”剩下的交给 CAP。我当初理解这个事情是拿快递类比你在网上下单商家先在自己的系统里录入订单业务表同时会打印一张面单消息表。面单打印成功这笔订单才能算接受。然后快递公司后台进程把面单拿走发布到 MQ即使快递车坏了面单还在可以重新叫车。相比直接让商家出门跑一趟同步调用面单模式显然更不容易丢。1.2 和 2PC、Saga、裸 MQ 的真实差距很多团队在选择方案时会在 2PC、Saga、直接发 MQ 之间纠结我把它们的本质差距列一下方案一致性业务侵入可用性主要代价2PC 分布式事务强一致高一般协调者单点事务期间锁资源吞吐受最慢节点制约Saga 补偿最终一致高高每个节点都要写补偿逻辑补偿漏写就是隐性事故裸 MQ先提交业务再发消息不一致低高业务成功后宕机消息就丢了业务失败但消息已发出现假事件CAP 本地消息表最终一致中高需要消息表与消费端幂等但对业务代码侵入可控2PC 看着最“稳”但实际投入成本巨大裸 MQ 最简单可它把“业务成功”和“消息发送”拆成了两个独立动作中间那个时间窗口就是丢消息的来源。CAP 走的路线是不追求同生共死的强一致而是保证业务成功后事件一定被发出消费端通过幂等把重复处理挡住。对绝大多数 API 集成场景这个性价比是最高的。1.3 什么场景该上 CAP什么场景别硬上适合上 CAP 的场景订单创建成功后要通知用户、加积分、更新搜索索引、同步 ERP这些后续操作允许“晚几秒但要到达”有回调的异步导出/导入API 只负责受理后台结果通过消息回写需要把一个本地操作广播给多个订阅方例如商品上下架后同步多个系统。不适合硬上的场景资金类强一致场景例如两账户之间必须同时扣款/入账且不允许中间状态CAP 只做最终一致不能当强一致分布式事务用对瞬时一致性要求极高、必须实时看到两边结果的功能还是同步接口更直接。我个人实际体会是即便资金类场景很多团队最终也会用 “CAP 幂等表 每日对账脚本” 来落地因为 2PC 的高成本在业务量起来以后很难承受。所以 CAP 的适用面比想象中更宽关键是你要接受“最终一致窗口”的存在并用对账和监控把它兜住。2. .NET 9 环境准备包怎么选、队列怎么挑、最小配置怎么写2.1 用 NuGet 把四个角色凑齐在 .NET 9 的 ASP.NET Core 项目里CAP 通过 NuGet 包引入按角色分四类DotNetCore.CAP核心包包含事务拦截、发布订阅、Dashboard 等队列包DotNetCore.CAP.RabbitMQ、DotNetCore.CAP.Kafka、DotNetCore.CAP.AzureServiceBus按实际中间件选一个数据库包DotNetCore.CAP.MySqlConnector、DotNetCore.CAP.SqlServer、DotNetCore.CAP.PostgreSQL要和业务库对应可选DotNetCore.CAP.InMemoryStorage、DotNetCore.CAP.InMemoryMessageQueue纯本地调试用。版本建议是直接安装 NuGet 上当前最新的 8.x 稳定版.NET 9 完全兼容。如果你是从 .NET 8 项目升级上来的CAP 包版本基本不用动只把目标框架改成 net9.0 即可。我最初就是因为没弄清楚这些包的角色一个项目里同时装了 RabbitMQ 和 Kafka 两个队列包虽然没出故障但纯属多余。2.2 RabbitMQ 还是 Kafka不只是看名气我大部分项目用 RabbitMQ理由很朴素运维成本低单机都能跑团队里没人学过消息中间件也能看 Web 管理界面队列、交换机模型和 CAP 的“订阅组”概念映射得最顺消息确认、死信、延迟重试都有现成能力。Kafka 什么时候更合适当业务事件量极高、需要按分区保证顺序、而且你已经有一套 Kafka 运维体系时CAP 接 Kafka 也支持。但注意Kafka 的分区顺序和消费者并发是相关的CAP 内部会尽量保证同组消费但你不应该依赖跨事件类型的全局顺序。数据库选择上建议直接用业务主库。因为本地消息表本来就要和业务表做同一个事务CAP 对 MySQL、SQL Server、PostgreSQL 的支持都很成熟。用 MySQL 8.x .NET 9 的 Pomelo EF Core 或 MySqlConnector我都在生产验证过。唯一要注意的是别用不支持数据库事务的存储来当消息表存储比如把事件记录写到 Redis 里那就彻底违背了本地消息表的设计初衷。2.3 最小启动配置跑通 Dashboard 才算第一步下面这段配置是我在 .NET 9 空模板项目上直接跑通的没有额外调整你可以先照抄var builder WebApplication.CreateBuilder(args); builder.Services.AddControllers(); builder.Services.AddDbContextAppDbContext(options options.UseMySql(builder.Configuration.GetConnectionString(Default), ServerVersion.AutoDetect(builder.Configuration.GetConnectionString(Default)))); builder.Services.AddCap(options { options.UseMySql(builder.Configuration.GetConnectionString(Default)); options.UseRabbitMQ(mq { mq.HostName builder.Configuration[RabbitMQ:Host]; mq.UserName builder.Configuration[RabbitMQ:User]; mq.Password builder.Configuration[RabbitMQ:Pass]; }); options.DefaultGroupName order.api.group; options.FailedRetryCount 5; options.FailedRetryInterval 30; options.SucceedMessageExpiredAfter TimeSpan.FromHours(24); options.FailedMessageExpiredAfter TimeSpan.FromDays(15); options.UseDashboard(); }); var app builder.Build(); app.MapControllers(); app.UseCapDashboard(); app.Run();逐项说明一下UseMySql让 CAP 在业务库里自动建消息表不需要手工迁移UseRabbitMQ指定 broker 连接DefaultGroupName定义默认消费者组同一组内消息会负载均衡不同组是广播后面细讲FailedRetryCount/FailedRetryInterval控制投递失败后的重试次数与间隔两个 ExpiredAfter成功消息和失败消息的保留时间到点自动清理防止cap.published、cap.received表无限膨胀Dashboard 需要同时有UseDashboard()配置和app.UseCapDashboard()中间件缺一个都看不到页面。CAP 启动后会自动建这几张表cap.published、cap.received、cap.lock。这件事很多人不知道以为要手工建表。不用框架启动时就自动完成了。2.4 消费者在哪里被发现的CAP 扫描的是程序集里带有[CapSubscribe]标记的 public 方法。所以你不需要在启动时显式注册某个消息路由只要消费者类注册到了 DICAP 就能拿到它。这里有个非常常见的坑消费者类忘记注册 DI运行时直接报错。后面第 4 节我会专门讲。3. 从订单 API 到消息消费一段完整的 CAP 集成链路3.1 控制器里的发布代码先定义事件 DTO。事件 DTO 和数据库实体要分开这个原则后面讲先按最佳实践写public class OrderCreatedEvent { public Guid OrderId { get; set; } public string UserId { get; set; } public decimal Amount { get; set; } public DateTime OccurredAt { get; set; } }控制器里创建订单的接口长这样[ApiController] [Route(api/orders)] public class OrdersController : ControllerBase { private readonly AppDbContext _db; private readonly ICapPublisher _cap; public OrdersController(AppDbContext db, ICapPublisher cap) { _db db; _cap cap; } [HttpPost] public async TaskIActionResult CreateOrder(CreateOrderRequest request) { using var transaction await _db.Database.BeginTransactionAsync(); var order new Order { Id Guid.NewGuid(), UserId request.UserId, Amount request.Amount, Status Created }; _db.Orders.Add(order); await _db.SaveChangesAsync(); await _cap.PublishAsync(order.created, new OrderCreatedEvent { OrderId order.Id, UserId request.UserId, Amount request.Amount, OccurredAt DateTime.UtcNow }); await transaction.CommitAsync(); return Ok(order.Id); } }为什么顺序是 SaveChanges - PublishAsync - Commit关键在于PublishAsync是在当前数据库事务里向cap.published插入一条消息记录。只要它在这个事务提交前被调用业务数据写入和消息记录就处于“同生共死”状态如果CommitAsync失败业务数据回滚消息记录也不存在如果CommitAsync成功但消息还没发布CAP 后台进程会从cap.published里捞出来补投所以不会丢。很多人容易写错的顺序是先 Commit 再 Publish。那样消息没进事务发布失败就只能自己补发CAP 的可靠性等于白装。这个顺序问题是我见过最普遍的误用。3.2 数据一致性背后的完整执行链路CAP 从发布到消费完整链路是这样的请求进入 Controller同一个数据库事务内orders 表插入 cap.published 插入一条 Pending 事件事务 CommitCAP 的后台 Collector 默认约 1 秒轮询一次 cap.published取出 Pending 消息将消息体序列化后发布到 RabbitMQ 的交换机收到 Broker ACK 后把状态标记为 Succeeded下游消费者从自己的队列拿到消息执行业务逻辑处理成功后CAP 向 MQ 返回 ACK并在本地 cap.received 记录已经消费若处理异常消息不 ACKCAP 按重试配置再次投递。这里必须反复强调一点CAP 提供的是“至少一次投递”不是“恰好一次投递”。也就是说极端情况下同一个事件可能被消费者处理两遍。所以订阅方必须有幂等设计这是使用 CAP 的底线要求。幂等的实现方式有很多最省事的是用业务唯一键判断。比如 Notify 场景里用订单号作为处理状态的唯一键[CapSubscribe(order.created)] public async Task OnOrderCreated(OrderCreatedEvent data) { var exists await _db.OrderNotifications .AnyAsync(n n.OrderId data.OrderId); if (exists) { return; // 已处理过直接跳过 } await _downstreamApi.NotifyOrderCreated(data); _db.OrderNotifications.Add(new OrderNotification { OrderId data.OrderId, CreatedAt DateTime.UtcNow }); await _db.SaveChangesAsync(); }这段代码的核心逻辑是先查是否处理过再调下游 API最后写处理记录。注意如果你把“调下游 API”放在“写处理记录”之前下游成功但本地写记录失败消息重投后仍然会重复调下游。所以更稳的做法是把处理记录和下游结果一起落库或者接受偶尔重复并让下游也支持幂等键。3.3 订阅侧调用下游 API 的正确姿势订阅消费者通常不只是更新本地库还要调用其他服务。这部分最容易翻车很多同事一上来就new HttpClient()。我强烈建议用IHttpClientFactory在 .NET 9 里还能配上原生 Resilience 管道public class OrderEventSubscriber { private readonly IHttpClientFactory _httpClientFactory; private readonly IConfiguration _config; public OrderEventSubscriber(IHttpClientFactory httpClientFactory, IConfiguration config) { _httpClientFactory httpClientFactory; _config config; } [CapSubscribe(order.created)] public async Task HandleOrderCreated(OrderCreatedEvent eventData) { var client _httpClientFactory.CreateClient(downstream); var payload new { eventData.OrderId, eventData.UserId, eventData.Amount }; var response await client.PostAsJsonAsync(/api/notify, payload, CancellationToken.None); if (!response.IsSuccessStatusCode) { // 401密钥问题不要自动重试500/502/504可以交给 CAP 重投 throw new InvalidOperationException($Downstream API returned {response.StatusCode}); } } }这里有个关键设计订阅函数内部不自己写 while 重试而是把异常抛出去让 CAP 负责重投。因为 CAP 重投附带持久化和 Dashboard 可视化出了问题你还能在 UI 上定位、手动重放而 HttpClient 内部重试对运维是不可见的。但也不是所有异常都适合让 CAP 重试。比如下游 API 返回 401 “incorrect api key provided”根本不是瞬时故障重试多少次都是 401。这时候应该记录告警日志把消息快速引入失败流程而不是让 CAP 一遍遍地把请求打到别人的网关上等于拿消息队列当暴力破解器。3.4 用 Dashboard 验证整条链路启动后访问/cap能看到 Succeeded、Failed、Queued、Received 四类消息统计以及每条消息的创建时间、重试次数、最后异常信息。对 Failed 消息Dashboard 可以直接点击重放把旧消息重新投递。还能看到消费者组和队列情况。我每次跑通新接口后习惯先去 Dashboard 确认这条order.created是 Succeeded再去消费者日志看是否收到。链路长了以后Dashboard 能帮你在“消息丢了”这种问题上省至少半小时排查时间。4. 发布与订阅的代码细节事件命名、消费者注册、序列化这些坑4.1 事件命名决定你的运维体验事件名是跨服务契约我坚持用“名词.过去分词”风格order.created、order.paid、user.registered、inventory.reserved。不要用OrderCreatedEvent这个类名当作路由类名一变所有订阅方都要跟着改而字符串事件名可以保持稳定。消息体建议使用独立 Event DTO不直接复用 EF 实体。原因很简单EF 实体往往包含导航属性、行版本、审计字段序列化出来内容很杂订阅方如果按契约反序列化新增字段不一定兼容。而 Event DTO 只暴露最小字段演进空间最大。我见过有项目直接把整个Order实体丢进消息里后来订单实体加了TenantId、导航对象消息体膨胀了三四倍订阅方解析还出了一堆反序列化异常。4.2 订阅类注册到 DI 的生命周期陷阱订阅类必须注册到 DI。比如上面那个OrderEventSubscriber需要显式加上builder.Services.AddScopedOrderEventSubscriber();因为消费者里通常会用到DbContext或IHttpClientFactory我建议注册成Scoped而不是Singleton。CAP 在消费时会从当前 scope 解析消费者实例Scoped 配合 AppDbContext 不会出现跨请求状态串互相污染的问题。常见报错信息是The type OrderEventSubscriber must be registered in DI container。看到这句话先检查 AddScoped/AddTransient 有没有写。还有同事为了省事把订阅方法写在 Controller 里虽然 CAP 也能扫到带[CapSubscribe]的 public 方法但不推荐因为 Controller 生命周期和 HTTP 请求绑定不适合做后台消费者。独立类最干净。4.3 消息序列化默认 JSON 够用但注意版本兼容CAP 默认将事件体序列化为 JSON 存储并投递。事件对象建议遵循几条规则字段类型稳定不要用dynamic或object传任意类型新版本可以加字段不要随便删字段因为可能还有旧消息没消费完反序列化侧对未知字段要宽容默认 JSON 反序列化通常没问题但不要开启严格未知字段报错模式尽量避免多态。消息体里一旦出现多态类型元信息类库升级后类型全名变化旧消息可能反序列化失败。我在实际项目里见过有人在事件里传JObject订阅方再ToObjectT()。确实能跑通但契约变得很模糊下游不知道消息结构到底什么样只能对着线上日志猜。不如显式定义事件类型让代码本身成为文档。4.4 消费失败后的重试两步重试容易叠加放大这是很多人忽略的CAP 有失败重投你在消费者内部调用 HTTP API 时如果再用重试策略两层重试叠起来下游一旦持续故障请求量会被放大好几倍。我实际遇到过下游数据库抖动 3 分钟消费者内部重试 3 次 CAP 重投 5 次瞬间产生大量并发请求直接把连接池打满。建议二选一方案 ACAP 负责重试订阅函数内部只调用下游一次失败即抛异常方案 BHttpClient 层做小幅重试只对超时和 5xx 重试 2 次CAP 只保留 1~2 次重投并且把FailedRetryInterval调大。我个人偏向方案 A。因为 CAP 重试带了持久化、Dashboard 可视化和手动重放能力比 HttpClient 内部重试更可控。如果选了方案 B务必在日志里打清楚“当前是第几次尝试”否则排查问题时根本分不清是哪种重试在打下游。另外FailedRetryInterval别设太短。默认 60 秒是比较合理的。设成 5 秒一旦 broker 或数据库抖动重试风暴会拖垮周边服务。这个参数和 HttpClient 的Timeout一样初期调小总觉得“恢复快”生产出事时才发现它是放大器。5. 我在生产环境里踩过的 CAP 相关的坑5.1 Dashboard 裸奔没有认证的运维后台最危险Dashboard 好用是真的好用裸奔也是真的裸奔。UseDashboard()不带任何身份认证知道/cap路径的人就能看到全部消息内容。消息体里如果有用户手机号、地址、订单金额等于给所有能看到网络的人开了一扇门。我的做法是生产环境通过反向代理只允许内网 IP 访问/cap或者至少套一层 Basic Auth / OAuth 前置认证服务如果直接暴露公网强烈建议彻底关闭 Dashboard改用日志和 metrics 做监控。nginx 配置示意location /cap { allow 10.0.0.0/8; deny all; proxy_pass http://127.0.0.1:8080; }看起来只是几行配置但很多人上线时忘了加等出了数据泄露事故再补就晚了。5.2 消息重试风暴问题不解决重试就是空转CAP 默认会把失败消息保留下来并反复重试直到超过FailedRetryCount。有一次我们的下游 API Key 被管理员重置了消费者收到 401 后还在按 CAP 配置重试消息队列被一堆注定失败的消息占满正常消息反而积压。这时候我的排查链路是先在下游 API 那边确认 Key 是否真的失效临时把消费者改成快速失败 告警避免无意义重试等 Key 恢复后在 Dashboard 上把积压的 Failed 消息手动重放重放前确认消费端幂等逻辑没问题因为历史消息里可能有部分已经处理成功。这个案例说明一个道理CAP 不会自动跳过坏消息你必须主动介入。所以告警机制比想象中重要。FailedThresholdCallback就是干这个的可以在里面接飞书、钉钉、企业微信机器人options.FailedThresholdCallback failed { // failed.Message 里有事件名、异常信息 // 在这里发送告警并带上消息 Id 方便后续重放 _alarmNotifier.NotifyAsync($CAP 消息进入死信: {failed.Message}); };5.3 密钥管理的实践不要在日志里出现完整的 API Key每次看到incorrect api key provided: sk-svcac****这种日志我都想再强调一次密钥管理。很多小团队为了方便把第三方 API Key 直接写进 appsettings.json 提交到 Git 仓库。等第三方网关升级鉴权策略后线上服务全部开始 401排查半天才发现是 Key 泄露被对方吊销了。我的标准流程本地开发用dotnet user-secrets保存密钥CI/CD 环境用环境变量或 Secret 注入生产容器用 Kubernetes Secret 挂载成环境变量或者统一引用密钥管理服务代码里只通过IConfiguration读一个入口。日志打印方面我建议统一封装一个脱敏逻辑记录 API Key 的 SHA-256 前 8 位和后 4 位例如key#sha256:abcd1234。这样排查 401 时能和第三方平台核对又不会完整泄露。5.4 .NET 9 周边 API 的兼容提示如果你在 .NET 9 项目里同时使用Microsoft.Extensions.Http.Resilience注意它的默认重试策略只针对 5xx 和超时不会自动跳过 401。所以我一般在配置里显式声明builder.Services.AddHttpClient(downstream, client { client.BaseAddress new Uri(builder.Configuration[Downstream:BaseUrl]!); client.DefaultRequestHeaders.Add(Authorization, $Bearer {apiKey}); }) .AddResilienceHandler(downstream-retry, r { r.AddRetry(new RetryStrategyOptions { MaxRetryAttempts 2, Delay TimeSpan.FromSeconds(2), ShouldHandle args args.Outcome.Result.StatusCode is HttpStatusCode.InternalServerError or HttpStatusCode.RequestTimeout ? ValueTask.FromResult(true) : ValueTask.FromResult(false) }); });这样 401 不会进入 HttpClient 重试而是快速抛给 CAP 的消费重投逻辑同时 CAP 测重试次数配小一点或者收到 401 时快速失败并触发告警。两个层面的重试职责就分清楚了。另外.NET 9 的静态分析规则比 .NET 8 更严格用PostAsyncStringContent时偶尔会收到 CA2264 之类的提示。建议直接改用PostAsJsonAsync或JsonContent.Create顺手把 Content-Type 正确性也解决了。这篇文章里的大部分内容是我把 CAP 集成进 .NET 9 在线订单 API 之后沉淀下来的。最深的体会可能有点反直觉CAP 这类框架引入之后真正决定系统稳定性的反而变成了消费端的幂等和密钥管理。如果你正在评估或者准备落地不要急着把所有业务都迁到消息上挑一个订单创建场景先跑通 Dashboard、重试、告警一个链路用顺了再铺开。后续消息量上来还可以考虑从 RabbitMQ 平移到 Kafka或者把 cap.published、cap.received 表做分区归档。但万变不离其宗本地消息表保证不丢消费端幂等保证不重这两条立住了分布式下的 API 集成才能睡得着觉。