- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
本文基于当前仓库(The NCC / CAP)中 PostgreSQL 存储官方文档 编写。CAP 是一款基于最终一致性思想、采用 Outbox 模式的微服务分布式事务与事件总线框架,PostgreSQL 是其官方全面支持的存储后端之一。读完本文,你将掌握如何为 CAP 引入 PostgreSQL 存储、理解
PostgreSqlOptions各配置项的底层含义,并学会通过 ADO.NET 原生事务与 Entity Framework Core 事务两种方式,把业务数据库操作与消息发布放进同一本地事务,从而保证"业务数据落库"与"事件消息可靠发布"的原子一致。
一、PostgreSQL 存储定位:CAP 的 Outbox 载体
CAP 的核心模型是在业务数据库的同一本地事务中写入业务数据与待发布消息(Outbox 模式),再由后台处理器将这些消息可靠地投递到消息队列。这意味着 CAP 必须拥有一个与业务库同源的数据存储,PostgreSQL 正是该存储的官方支持选项之一。
从仓库源码看,PostgreSQL 存储模块 src/DotNetCore.CAP.PostgreSql 主要承担三件事:
- 建表与初始化:由
PostgreSqlStorageInitializer在启动时自动创建published、received(以及可选lock)三张表,参见 IStorageInitializer.PostgreSql.cs; - 消息持久化读写:由
PostgreSqlDataStorage实现IDataStorage接口,负责消息的存储、状态流转、重试与加锁,参见 IDataStorage.PostgreSql.cs; - 本地事务集成:由
PostgreSqlCapTransaction与CapTransactionExtensions提供将 CAP 发布动作挂载到业务事务的能力,参见 ICapTransaction.PostgreSql.cs。
二、安装与基本配置
2.1 安装 NuGet 包
在包管理器控制台中执行:
PM> Install-Package DotNetCore.CAP.PostgreSql从 DotNetCore.CAP.PostgreSql.csproj 可以确认,该包当前面向net8.0;net9.0;net10.0三个目标框架,并依赖Npgsql(版本 10.0.2)与对应版本的Microsoft.EntityFrameworkCore.Relational。这意味着使用 PostgreSQL 存储时,请确保项目目标框架与上述范围匹配。
2.2 在 ConfigureServices 中注册 CAP
在Startup.cs(或最小托管模型的Program.cs)的ConfigureServices方法中添加:
public void ConfigureServices(IServiceCollection services) { // ... services.AddCap(x => { x.UsePostgreSql(opt => { // PostgreSqlOptions 配置项 }); // x.UseXXX ... (消息队列传输配置,例如 Kafka、RabbitMQ 等) }); }仓库示例 Sample.Kafka.PostgreSql/Program.cs 给出了一个完整的最小托管示例:它先通过UseNpgsql注册业务DbContext,随后在AddCap中同时调用x.UsePostgreSql(...)与x.UseKafka(...),并注册 Dashboard,可以直接作为集成参考。
关于连接字符串写法,示例中给出的是:
public const string DbConnectionString = "User ID=postgres;Password=mysecretpassword;Host=127.0.0.1;Port=5432;Database=postgres;";2.3 PostgreSqlOptions 配置项详解
官方文档给出如下参数表:
| NAME | DESCRIPTION | TYPE | DEFAULT |
|---|---|---|---|
| Schema | 数据库 schema | string | cap |
| ConnectionString | 数据库连接字符串 | string | 无 |
| DataSource | Npgsql 数据源 | NpgsqlDataSource | 无 |
结合源码可以进一步理解这三项的实际行为:
- Schema:默认值为
cap(常量EFOptions.DefaultSchema,见 CAP.EFOptions.cs)。存储初始化器会执行CREATE SCHEMA IF NOT EXISTS "cap"并以"schema"."published"、"schema"."received"、"schema"."lock"的形式引用表,参见 IStorageInitializer.PostgreSql.cs。 - ConnectionString与DataSource:二者二选一即可。
PostgreSqlOptions.CreateConnection()的优先级是:若配置了DataSource则调用DataSource.CreateConnection(),否则用ConnectionString新建NpgsqlConnection,参见 CAP.PostgreSqlOptions.cs。
UsePostgreSql扩展方法提供两种重载(见 CAP.Options.Extensions.cs):
// 方式一:直接传连接字符串 x.UsePostgreSql("User ID=postgres;Password=...;Host=127.0.0.1;Port=5432;Database=cap;"); // 方式二:通过委托精细配置 x.UsePostgreSql(opt => { opt.Schema = "cap"; opt.ConnectionString = "User ID=postgres;Password=...;Host=127.0.0.1;Port=5432;Database=cap;"; });注册时内部会通过PostgreSqlCapOptionsExtension完成三件事(见 CAP.PostgreSqlCapOptionsExtension.cs):
- 注册存储标识
CapStorageMarkerService("PostgreSql"),用于多存储扩展共存时的标识与校验; - 注册
IConfigureOptions<PostgreSqlOptions>配置器; - 将
PostgreSqlDataStorage(实现IDataStorage)与PostgreSqlStorageInitializer(实现IStorageInitializer)注册为单例。
此外,还有一种基于 Entity Framework Core 的注册方式UseEntityFramework<TContext>()(同见 CAP.Options.Extensions.cs):此时 CAP 会从TContext的IDbContextOptions扩展中自动提取DataSource或ConnectionString,无需重复配置连接字符串。ConfigurePostgreSqlOptions中还对"在 DbContext 内注入ICapPublisher"的情况做了循环引用保护(会抛出明确异常,提示改用x.UsePostgreSql()直接配置存储),参见 CAP.PostgreSqlOptions.cs。
三、自动建表:CAP 在 PostgreSQL 中创建的库表结构
首次启动时,PostgreSqlStorageInitializer.InitializeAsync会执行建库脚本(见 IStorageInitializer.PostgreSql.cs),核心结构如下:
- received 表:
Id(BIGINT 主键)、Version、Name、Group、Content(TEXT)、Retries、Added、ExpiresAt、StatusName,并建有(ExpiresAt, StatusName)与(Version, ExpiresAt, StatusName)两个索引; - published 表:结构类似 received(无
Group列),同样建有上述两个索引; - lock 表(仅在
CapOptions.UseStorageLock开启时创建):包含Key、Instance、LastLockTime三列,初始化时通过ON CONFLICT DO NOTHING幂等地插入publish_retry_{Version}与received_retry_{Version}两把锁记录。
这些索引与锁表正是 CAP 重试处理器、过期清理与多实例抢占调度得以高效运行的基础。所有 DDL 均使用IF NOT EXISTS,可安全地在多次启动中重复执行。
四、本地事务发布:保证业务与消息的原子性
CAP 的核心用法是"事务内发布":在业务数据库的本地事务中同时写入业务数据和待发布消息,二者要么一起成功、要么一起回滚,从而避免"业务成功但消息丢失"或"消息发出但业务回滚"的不一致问题。
事务发布依赖两个先决条件:
- 通过
services.AddCap(x => { x.UsePostgreSql(...); ... })注册了 CAP 与 PostgreSQL 存储; - 在业务代码中注入
ICapPublisher _capBus。
4.1 ADO.NET + Npgsql 原生事务
private readonly ICapPublisher _capBus; using (var connection = new NpgsqlConnection("ConnectionString")) { using (var transaction = connection.BeginTransaction(_capBus, autoCommit: false)) { // 你的业务代码 connection.Execute("insert into test(name) values('test')", transaction: (IDbTransaction)transaction.DbTransaction); _capBus.Publish("sample.rabbitmq.mysql", DateTime.Now); transaction.Commit(); } }说明:
connection.BeginTransaction(_capBus, autoCommit: false)是IDbConnection上的扩展方法,它会创建PostgreSqlCapTransaction实例并挂到publisher.Transaction上(见 ICapTransaction.PostgreSql.cs);_capBus.Publish(...)将消息以 Outbox 形式写入published表,此时并未真正发出队列消息;- 调用
transaction.Commit()后,PostgreSqlCapTransaction.Commit()会先提交数据库事务,再通过Flush()把消息派发给传输层(见 ICapTransaction.PostgreSql.cs),从而保证"先落库、后投递"的最终一致语义; - 若业务异常导致回滚,消息记录也会一并回滚,不会出现幽灵消息。
扩展方法还提供了带隔离级别与异步的重载:BeginTransaction(IsolationLevel, publisher, autoCommit)、BeginTransactionAsync(...),以及通过transaction.DbTransaction访问底层IDbTransaction的能力。autoCommit: true时,Publish会立即自动提交事务;而false时则必须显式Commit()。
4.2 Entity Framework Core 事务
private readonly ICapPublisher _capBus; using (var trans = dbContext.Database.BeginTransaction(_capBus, autoCommit: false)) { dbContext.Persons.Add(new Person() { Name = "ef.transaction" }); _capBus.Publish("sample.rabbitmq.mysql", DateTime.Now); dbContext.SaveChanges(); trans.Commit(); }说明:
dbContext.Database.BeginTransaction(_capBus, autoCommit: false)是DatabaseFacade上的扩展方法,返回的trans实际是包装了PostgreSqlCapTransaction的CapEFDbTransaction(见 IDbContextTransaction.CAP.cs),因此它同时实现了IDbContextTransaction与IInfrastructure<DbTransaction>,与 EF Core 的事务语义完全兼容;- 与 ADO.NET 方式一致,
Commit()内部会先提交 EF 事务,再调用Flush()投递消息; - 推荐顺序是:先执行业务数据变更(此处示例为
dbContext.Persons.Add(...)),再Publish,最后SaveChanges()与Commit()。
DatabaseFacade同样支持带隔离级别的同步/异步重载:BeginTransaction(isolationLevel, publisher, autoCommit)与BeginTransactionAsync(...),参见 ICapTransaction.PostgreSql.cs。
4.3 两种方式的差异与选型
| 对比维度 | ADO.NET 方式 | EF Core 方式 |
|---|---|---|
| 适用场景 | 使用 Dapper、原生 SQL 或未引入 EF 的项目 | 已使用 Entity Framework Core 的项目 |
| 入口 API | IDbConnection.BeginTransaction(publisher, autoCommit) | DatabaseFacade.BeginTransaction(publisher, autoCommit) |
| 底层事务类型 | IDbTransaction/DbTransaction | IDbContextTransaction(经CapEFDbTransaction包装) |
| 连接配置 | 需自行提供NpgsqlConnection连接字符串 | 直接复用DbContext的 Npgsql 连接 |
五、总结与注意事项
- 包版本匹配:
DotNetCore.CAP.PostgreSql面向net8.0/net9.0/net10.0,基于Npgsql 10.0.2,引用前请核对目标框架与 Npgsql 版本兼容性(见 DotNetCore.CAP.PostgreSql.csproj); - Schema 可自定义:默认建在
capschema 下,业务上若需多租户或隔离,可修改opt.Schema,表名与索引会随 schema 自动调整; - ConnectionString 与 DataSource 二选一:两者都提供时优先使用
DataSource(CAP.PostgreSqlOptions.cs),若使用UseEntityFramework<TContext>注册则可免去连接配置; - 事务内发布是保证一致性的关键:务必在
Commit()之前完成Publish,且不要跨事务复用ICapPublisher的Transaction状态;autoCommit: false时必须显式提交,否则消息不会投递; - 多实例部署:若开启
UseStorageLock,CAP 会在lock表中通过AcquireLockAsync/RenewLockAsync/ReleaseLockAsync实现分布式任务抢占(见 IDataStorage.PostgreSql.cs),保证重试与清扫处理器在多节点下只由一个实例执行。
至此,你已经完成了 CAP + PostgreSQL 存储从安装、配置到事务化发布的全链路搭建,可以在此基础上继续配置你选择的传输层(如 Kafka、RabbitMQ、Azure Service Bus 等),构建具备最终一致性的微服务事件总线。
- 后端
- 消息队列
- 微服务
【免费下载链接】CAP
基于最终一致性的微服务分布式事务解决方案,也是一种采用 Outbox 模式的事件总线。
相关推荐
使用 CAP 与 MongoDB 构建基于 Outbox 模式的分布式消息存储与本地事务
使用 CAP 与 MongoDB 构建基于 Outbox 模式的分布式消息存储与本地事务 本指南系统讲解如何将 CAP(.NET 微服务分布式事务与事件总线框架
后端消息队列微服务消息路由CAP 使用 MongoDB 作为消息存储:配置、事务发布与源码级原理剖析
CAP 使用 MongoDB 作为消息存储:配置、事务发布与源码级原理剖析 本文以 CAP(基于 Outbox 模式的 .NET 分布式事务解决方案)的官方英文
后端消息队列微服务CAP 事件总线 SQL Server 存储接入指南:配置、Outbox 表结构与本地消息事务实战
CAP 事件总线 SQL Server 存储接入指南:配置、Outbox 表结构与本地消息事务实战 导读 本文是 DotNetCore.CAP 分布式事务框架使
后端消息队列微服务消息路由
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考