
后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载本文基于当前仓库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 IDpostgres;Passwordmysecretpassword;Host127.0.0.1;Port5432;Databasepostgres;;2.3 PostgreSqlOptions 配置项详解官方文档给出如下参数表NAMEDESCRIPTIONTYPEDEFAULTSchema数据库 schemastringcapConnectionString数据库连接字符串string无DataSourceNpgsql 数据源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 IDpostgres;Password...;Host127.0.0.1;Port5432;Databasecap;); // 方式二通过委托精细配置 x.UsePostgreSql(opt { opt.Schema cap; opt.ConnectionString User IDpostgres;Password...;Host127.0.0.1;Port5432;Databasecap;; });注册时内部会通过PostgreSqlCapOptionsExtension完成三件事见 CAP.PostgreSqlCapOptionsExtension.cs注册存储标识CapStorageMarkerService(PostgreSql)用于多存储扩展共存时的标识与校验注册IConfigureOptionsPostgreSqlOptions配置器将PostgreSqlDataStorage实现IDataStorage与PostgreSqlStorageInitializer实现IStorageInitializer注册为单例。此外还有一种基于 Entity Framework Core 的注册方式UseEntityFrameworkTContext()同见 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 表IdBIGINT 主键、Version、Name、Group、ContentTEXT、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与IInfrastructureDbTransaction与 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 的项目入口 APIIDbConnection.BeginTransaction(publisher, autoCommit)DatabaseFacade.BeginTransaction(publisher, autoCommit)底层事务类型IDbTransaction/DbTransactionIDbContextTransaction经CapEFDbTransaction包装连接配置需自行提供NpgsqlConnection连接字符串直接复用DbContext的 Npgsql 连接五、总结与注意事项包版本匹配DotNetCore.CAP.PostgreSql面向net8.0/net9.0/net10.0基于Npgsql 10.0.2引用前请核对目标框架与 Npgsql 版本兼容性见 DotNetCore.CAP.PostgreSql.csprojSchema 可自定义默认建在capschema 下业务上若需多租户或隔离可修改opt.Schema表名与索引会随 schema 自动调整ConnectionString 与 DataSource 二选一两者都提供时优先使用DataSourceCAP.PostgreSqlOptions.cs若使用UseEntityFrameworkTContext注册则可免去连接配置事务内发布是保证一致性的关键务必在Commit()之前完成Publish且不要跨事务复用ICapPublisher的Transaction状态autoCommit: false时必须显式提交否则消息不会投递多实例部署若开启UseStorageLockCAP 会在lock表中通过AcquireLockAsync/RenewLockAsync/ReleaseLockAsync实现分布式任务抢占见 IDataStorage.PostgreSql.cs保证重试与清扫处理器在多节点下只由一个实例执行。至此你已经完成了 CAP PostgreSQL 存储从安装、配置到事务化发布的全链路搭建可以在此基础上继续配置你选择的传输层如 Kafka、RabbitMQ、Azure Service Bus 等构建具备最终一致性的微服务事件总线。赞分享后端消息队列微服务【免费下载链接】CAP基于最终一致性的微服务分布式事务解决方案也是一种采用 Outbox 模式的事件总线。项目地址https://gitcode.com/dotnetcore/CAP点击查看免费下载相关推荐使用 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 分布式事务框架使后端消息队列微服务消息路由上一篇3秒搞定网页图片格式转换Save Image as Type Chrome扩展终极指南下一篇3秒搞定图片格式转换Chrome扩展神器Save Image as Type使用指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考