
后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载导读本文围绕 DotNetCore.CAP基于 Outbox 模式的分布式事务与事件总线解决方案如何接入 MySQL 存储展开完整讲解 NuGet 包安装、UseMySql配置项、MySqlOptions参数含义以及 ADO.NET / Entity Framework Core 两种事务内发布消息的标准写法并结合仓库源码与官方示例Sample.RabbitMQ.MySql说明底层表结构、消息状态流转与事务提交时消息真正落库发送的机制。读完本文你将能独立完成MySQL 存储 本地数据库事务 消息可靠发布这一经典组合的落地并理解其为什么能保证业务数据与消息的一致性。一、为什么选择 MySQL 作为 CAP 的存储CAP 是分布式事务 事件总线的微服务解决方案核心思路是 Outbox发件箱模式业务数据与待发送消息写入同一个本地数据库事务随后后台组件把消息可靠地投递给消息队列RabbitMQ、Kafka、SQS 等消费端消费并记录结果最终达成最终一致性。因此存储层是 CAP 的基石——它既要保存待发布/已接收的消息又要承载业务事务。MySQL 作为开源关系型数据库与 CAP 有完整的集成支持CAP 官方为其提供了独立的 NuGet 包DotNetCore.CAP.MySql包含存储实现MySqlDataStorage、表结构初始化器MySqlStorageInitializer、监控 APIMySqlMonitoringApi以及事务扩展MySqlCapTransaction代码位于 src/DotNetCore.CAP.MySql。二、安装 NuGet 包在项目中使用 MySQL 存储首先通过 NuGet 安装对应包在 Visual Studio 的包管理器控制台中执行或在dotnet add package DotNetCore.CAP.MySql的命令行中执行PM Install-Package DotNetCore.CAP.MySql安装后包内的MySqlCapOptionsExtension会自动完成以下注册参见 CAP.MySqlCapOptionsExtension.cs注册存储标记服务CapStorageMarkerService(MySql)注册IDataStorage实现为MySqlDataStorage注册IStorageInitializer实现为MySqlStorageInitializer注册并配置MySqlOptions。三、配置 UseMySql在Startup.cs或 .NET 6 的Program.cs的ConfigureServices中添加 CAP 配置public void ConfigureServices(IServiceCollection services) { // ... services.AddCap(x { x.UseMySql(opt { // MySqlOptions opt.ConnectionString Server127.0.0.1;Databasecap;Uidroot;Pwdmy-secret-pw;; opt.TableNamePrefix cap; }); // x.UseXXX ... 例如消息队列传输配置 }); }仓库中的官方示例 Sample.RabbitMQ.MySql/Program.cs 采用字符串重载直接传入连接串同时组合了 RabbitMQ 传输与 Dashboard 监控面板builder.Services.AddCap(x { x.UseMySql(AppDbContext.ConnectionString); x.UseRabbitMQ(localhost); x.UseDashboard(); });MySqlOptions 配置项名称描述类型默认值TableNamePrefixCAP 表名前缀stringcapConnectionString数据库连接字符串stringnull对照源码 CAP.EFOptions.csMySqlOptions继承自EFOptions其中TableNamePrefix默认值为cap所有 CAP 表都以该前缀命名Version为数据版本标记内部使用默认v1写入消息表时用于区分不同版本的 CAP 数据。UseMySql提供了两种重载参见 CAP.Options.Extensions.cs// 方式一直接传连接字符串 options.UseMySql(Server127.0.0.1;Databasecap;Uidroot;Pwd...); // 方式二通过委托配置 MySqlOptions options.UseMySql(opt { opt.ConnectionString ...; opt.TableNamePrefix mycap; });若同时使用 Entity Framework Core还可以通过UseEntityFrameworkTContext()让 CAP 从已注册的 DbContext 中自动读取连接字符串见 CAP.Options.Extensions.cs 与 CAP.MySqlOptions.cs 中ConfigureMySqlOptions的自动获取逻辑。注意如果 DbContext 中注入了ICapPublisher源码会抛出异常提示改用UseMySql()直接配置连接串以避免循环引用。连接字符串建议示例 AppDbContext.cs 中的连接串格式为Server127.0.0.1;Databasecap;Uidroot;Pwdmy-secret-pw;实际项目中应使用 MySQL 官方驱动如 MySqlConnector支持的完整连接参数包括端口Port3306、字符集Charsetutf8mb4、连接池设置等并按环境通过配置中心或环境变量管理避免硬编码。四、CAP 在 MySQL 中自动创建的表结构CAP 首次启动时MySqlStorageInitializer.InitializeAsync会自动执行建表脚本见 IStorageInitializer.MySql.cs。表名由前缀派生{prefix}.published—— 已发布消息表{prefix}.received—— 已接收消息表{prefix}.lock—— 分布式锁表仅当UseStorageLock开启时创建。建表脚本见 IStorageInitializer.MySql.cs以 InnoDB 引擎、utf8mb4字符集创建关键字段包括Idbigint 主键、Versionvarchar版本标记、Name消息主题名、Contentlongtext序列化后的消息体、Retries重试次数、Added/ExpiresAt时间戳、StatusName状态Scheduled / Succeeded / Failed / Delayed 等received表额外包含Group消费组字段两表均建有IX_Version_ExpiresAt_StatusName与IX_ExpiresAt_StatusName索引支撑重试扫描与过期清理查询。此外初始化器会解析服务器版本ServerVersion.Parse以判断是否支持FOR UPDATE SKIP LOCKEDMySQL 8.0 或 MariaDB 10.6 返回true见 IStorageInitializer.MySql.cs该能力用于调度与清理时的行级跳过锁定提升多实例并发下的吞吐。五、事务内发布消息核心用法CAP 事务发布的核心思想业务 SQL 与消息写入StoreMessageAsync向published表插入必须在同一个数据库事务中完成只有事务提交后消息才会真正被派发Flush到消息队列。底层由MySqlCapTransaction在Commit时先提交数据库事务再调用Flush()见 ICapTransaction.MySql.cs 与 ICapTransaction.Base.cs。5.1 ADO.NETMySqlConnection / Dapper方式private readonly ICapPublisher _capBus; using (var connection new MySqlConnection(AppDbContext.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)是 CAP 提供的扩展方法见 ICapTransaction.MySql.cs它会将MySqlCapTransaction与当前连接绑定并写入publisher.Transactiontransaction.DbTransaction暴露底层IDbTransaction用于传递给 Dapper 等 ORM 执行业务 SQLautoCommit参数为true时发布消息即自动提交事务为false时必须手动调用Commit()。官方示例 ValuesController.cs 使用了await connection.BeginTransactionAsync(_capBus, true)的异步简化写法若业务代码抛出异常事务回滚Rollback()会丢弃缓冲的消息业务数据与消息都不会生效从而保证一致性。5.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, ...)是DatabaseFacade上的 CAP 扩展见 ICapTransaction.MySql.cs内部通过ActivatorUtilities创建MySqlCapTransaction并包装为CapEFDbTransaction见 IDbContextTransaction.CAP.cs因此trans.Commit()最终会调用 CAP 事务的Commit()推荐顺序先Publish再SaveChanges再Commit确保消息记录与业务实体处于同一事务也提供异步版本BeginTransactionAsync(_capBus, ...)供 async 场景使用。六、无事务发布与延迟消息补充场景在实际项目中并非所有消息都必须放在事务内。仓库示例 ValuesController.cs 还演示了两种常见用法// 无事务直接发布 await _capBus.PublishAsync(sample.rabbitmq.mysql, DateTime.Now, cancellationToken: HttpContext.RequestAborted); // 延迟发布指定秒数后投递 await _capBus.PublishDelayAsync(TimeSpan.FromSeconds(delaySeconds), sample.rabbitmq.test, $publish time:{DateTime.Now}, delay seconds:{delaySeconds});PublishDelayAsync会在消息头写入DelayTime标记MySqlCapTransaction.FlushAsync检测到延迟消息后将其交给调度器EnqueueToScheduler而非立即发布见 ICapTransaction.Base.cs随后由ScheduleMessagesOfDelayedAsync见 IDataStorage.MySql.cs按到期时间扫描并投递。七、存储层的可靠性机制源码视角结合 IDataStorage.MySql.cs 可以看出 MySQL 存储实现承担的关键职责消息入库StoreMessageAsync在事务内把消息插入published表初始状态为Scheduled失败消息则通过StoreReceivedExceptionMessageAsync记录并设初始重试次数为FailedRetryCount状态流转ChangePublishStateAsync/ChangeReceiveStateAsync随消息投递进度更新StatusNameScheduled → Succeeded / Failed支持在事务内以同一连接执行更新重试机制GetPublishedMessagesOfNeedRetry/GetReceivedMessagesOfNeedRetry定期扫描Retries FailedRetryCount且状态为 Failed / Scheduled 的消息配合 CAP 的FailedRetryCount、FailedMessageExpiredAfter选项实现补偿重试过期清理DeleteExpiresAsync批量删除已成功或已失败且超过过期时间的消息默认每批 1000 条MySQL 8 下使用FOR UPDATE SKIP LOCKED避免多实例重复清理分布式锁AcquireLockAsync/RenewLockAsync/ReleaseLockAsync基于lock表实现实例间互斥防止多节点同时调度同一批消息。八、完整的可运行示例仓库中 Sample.RabbitMQ.MySql 是MySQL 存储 RabbitMQ 传输的完整可运行示例目录结构包括Program.cs—— 服务注册与 CAP 配置UseMySqlUseRabbitMQUseDashboardAppDbContext.cs—— 连接字符串常量Controllers/ValuesController.cs—— 无事务发布、延迟发布、ADO.NET 事务发布以及[CapSubscribe]订阅者示例appsettings.json—— 应用配置。运行该示例前需准备 MySQL 实例与 RabbitMQ并确认连接串与传输配置正确启动后访问 Dashboard 可查看消息发布/接收状态与监控数据。九、注意事项与最佳实践事务与消息必须同库事务发布的前提是业务库与 CAP 消息库是同一个 MySQL 数据库同一个连接串跨库无法保证原子性连接字符串安全使用配置管理工具注入连接串生产环境禁用明文密码示例中的Pwdmy-secret-pw仅为本地演示字符集建表脚本默认utf8mb4业务库也应保持一致以支持完整 UnicodeMySQL 版本FOR UPDATE SKIP LOCKED仅在 MySQL 8.0 或 MariaDB 10.6 生效旧版本自动回退为普通FOR UPDATE不影响功能正确性异步优先ASP.NET Core 场景推荐使用BeginTransactionAsync/PublishAsync等异步 API避免阻塞线程池。结语通过DotNetCore.CAP.MySql包CAP 将 MySQL 变成了可靠的消息发件箱配置只需UseMySql一行接入事务发布则只需在 ADO.NET 或 EF Core 事务内调用_capBus.Publish配合自动建表、状态流转、重试与过期清理机制即可在微服务架构中以最终一致性完成分布式事务。结合 src/DotNetCore.CAP.MySql 源码与 Sample.RabbitMQ.MySql 示例你可以进一步定制表前缀、连接串、重试次数等参数快速落地生产环境。赞分享后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载相关推荐DotNetCore.CAP 集成 SQL Server 存储配置、建表机制与事务消息发布实战指南DotNetCore.CAP 集成 SQL Server 存储配置、建表机制与事务消息发布实战指南 导读 本文以 CAPDotNetCore.CAP事件总后端消息队列微服务CAP 使用 MongoDB 作为消息存储配置、事务发布与源码级原理剖析CAP 使用 MongoDB 作为消息存储配置、事务发布与源码级原理剖析 本文以 CAP基于 Outbox 模式的 .NET 分布式事务解决方案的官方英文后端消息队列微服务CAP Transports 完全指南为 DotNetCore.CAP 事件总线选型、配置并理解消息传输层CAP Transports 完全指南为 DotNetCore.CAP 事件总线选型、配置并理解消息传输层 本文以 docs/content/user gui后端消息队列微服务上一篇agent-plugins的Dart技能详解八dart-setup-ffi-assets FFI资源配置下一篇Ink Kit表单组件Input、Checkbox、Radio最佳实践创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考