BitzOrcas 使用 DotNetCore.CAP 连接 SQL Server Outbox 与 RabbitMQ。它解决“先提交数据库还是先发消息”的双写问题,但不会把分布式消息变成 exactly-once。生产者获得事务性 Outbox;消费者仍应按至少一次投递设计。
关键路径图
下图展示了 CAP 事务性发件箱(Transactional Outbox)从领域事件收集、数据库本地事务同提交到后台投递给消息中间件的完整拓扑流程。
生产路径
CapSqlSugarUnitOfWork 在共享 DbConnection 上开启 CAP 感知事务,并把同一 DbTransaction 注入 SqlSugar。业务写入和 CAP Outbox INSERT 随同一个事务提交或回滚。
Command ↓CAP-aware SQL transaction ├─ business rows └─ CAP Outbox row ↓ commitCAP dispatcher → RabbitMQ → consumer只有在这个事务仍然打开时调用 CAP 发布,才能获得原子双写保证。提交后的领域事件自动桥接属于 best-effort,边界见事件抽象层。
主题与契约
发布者拥有主题和载荷契约。强类型事件可以用 [IntegrationTopic("tickets.ticket.opened.v1")] // 主题常量带 .v1,与 EventCatalog 登记一致 固定主题;没有显式主题时,CAP 适配器回退到类型全名。长期对外主题应显式声明,并记录版本策略。
契约演进遵循以下原则:
- 优先增加可选字段,不删除或改变既有字段含义。
- 发布者不暴露数据库实体和内部枚举序号。
- 事件 ID、租户 ID 和发生时间是排错与幂等的基础。
- 破坏性变更使用新 topic 或新版本,并允许新旧消费者并行迁移。
订阅主题常量由消费者模块的 Contracts/Application 目录暴露给模块治理报告;运行时消费者在 Infrastructure 中实现 ICapSubscribe 和 [CapSubscribe]。声明报告与运行时属性必须保持一致。
消费者规则
消费者收到消息后先验证契约与租户,再执行 owner-local 用例。常见幂等方式:
- 以
eventId + consumer建唯一处理记录。 - 对目标聚合使用条件状态迁移。
- 投影表使用业务唯一键 Upsert。
成功写入业务效果和幂等标记应尽量放在同一事务。不要先标记“已处理”再执行副作用,也不要通过进程内 HashSet 去重。
消费失败应抛出异常,让 CAP 记录并重试。当前 CAP 配置的失败重试次数为 5。全局 CapConsumerAuditFilter 会产出后台执行审计;运维模块还能查询和受控重试 CAP 失败消息。
事务内发布用例
下面的用例把状态保存与事件发布放在同一个由事务行为包围的 Handler 中。TicketOpenedIntegrationEvent 是由 Tickets Contracts 拥有的稳定契约。
public sealed class OpenTicketHandler( ITicketStore tickets, IIntegrationEventPublisher<TicketOpenedIntegrationEvent> events){ public async Task<Result> Handle(OpenTicket command, CancellationToken ct) { // 聚合先保护状态迁移;预期失败直接返回,事务行为会回滚。 var ticket = await tickets.GetAsync(command.TicketId, ct); var opened = ticket.Open(command.Subject); if (opened.IsFailure) return opened;
// 保存与 CAP Publish 都发生在仍然打开的 CAP 感知事务中。 await tickets.SaveAsync(ticket, ct); await events.PublishAsync( TicketOpenedIntegrationEvent.From(ticket), ct);
return Result.Success(); }}Handler 不调用 CommitAsync();TransactionPipelineBehavior 根据 Result 与异常决定提交或回滚。如果把 PublishAsync 移到提交后回调,就失去本例依赖的原子 Outbox 语义。
幂等消费用例
public sealed class TicketOpenedConsumer( IProcessedEventStore processed, IReportingTicketProjection projection){ [CapSubscribe(TicketOpenedIntegrationEvent.Topic)] public async Task Consume(TicketOpenedIntegrationEvent message, CancellationToken ct) { // eventId + consumerName 必须有数据库唯一约束;进程内缓存不能替代它。 if (await processed.ExistsAsync(message.EventId, nameof(TicketOpenedConsumer), ct)) return;
// 投影更新与已处理标记应在消费者本地事务中一起提交。 await projection.UpsertAsync(message.TenantId, message.TicketId, message.Subject, ct); await processed.AddAsync(message.EventId, nameof(TicketOpenedConsumer), ct); }}并发收到同一事件时,唯一约束是最后防线。若 AddAsync 冲突,消费者应把“另一执行已完成”识别为幂等成功;若投影尚未提交,则继续按失败处理并让 CAP 重试。
故障边界
| 故障 | 预期行为 |
|---|---|
| 业务事务回滚 | 业务行与 Outbox 行一起回滚 |
| RabbitMQ 暂时不可用 | Outbox 保留,恢复后继续投递 |
| 消费者抛异常 | CAP 重试并记录失败 |
| 消费成功但确认丢失 | 可能重复投递,依赖业务幂等 |
| 提交后 best-effort 发布失败 | 业务已成功,只有告警,需单独补偿 |
上线门禁
- SQL Server、CAP 表和 RabbitMQ 配置已通过健康/就绪验证。
- 主题目录与运行时订阅一致,没有无人消费或拼写漂移。
- 每个消费者有重复消息、乱序、旧版本和毒消息测试。
- 失败消息积压、重试率、消费耗时和最老消息年龄有监控。
- 关键流程能从事件 ID、CorrelationId 追到生产者与消费者审计。
源码与验证入口
- CAP 感知工作单元:
CapSqlSugarUnitOfWork、CapEfCoreUnitOfWork。 - 类型化发布适配器:
CapIntegrationEventPublisher<TEvent>。 - 领域事件桥接:
DomainEventDispatchPipelineBehavior与IIntegrationTopicRegistry。 - 消费审计:
CapConsumerAuditFilter;失败消息运维位于 OpsExtension。 - 集成证据应覆盖业务行与 Outbox 同提交/同回滚、Broker 恢复后补投、重复消费和受控重试。
事件契约和端口的选择见事件抽象层。