CAP不是自动分布式事务,需数据库表+消息队列+后台调度器协同;PublishAsync仅写本地Cap.Published表并标记Processing,由Dispatcher轮询投递;未配连接串或事务未提交会导致消息丢失;订阅依赖特性标注与程序集扫描;幂等需业务层实现,推荐用MessageId查Cap.Received表去重;调度器卡住常因Bootstrapper ExecuteAsync阻塞或配置缺失。
CAP 不是“开启事务开关就能自动分布式”,它靠的是数据库表 + 消息队列 + 后台调度器三者协同,消息必须先落库、再发队列、再由消费者执行,漏掉任一环都会导致最终不一致。
ICapPublisher.PublishAsync 为什么没立刻发到 Kafka/RabbitMQ?
CAP 的
方法根本不会直连消息中间件。它只做两件事:把消息序列化后写入本地数据库的
表,并标记为
状态。真正的投递由后台
定期轮询该表触发。
如果你没配
或数据库连接失败,
会直接抛
:“Failed to insert message into database”
如果业务代码在
中调用
,但没显式
,那这条记录会随事务回滚——消息彻底消失,不会重试
支持传入
字典,但只有
、
等内置 key 会被调度器识别,自定义 header 需要在消费者端手动解析
Consumer 方法不触发?检查这三个地方
CAP 的订阅发现完全基于程序集扫描 + 特性标注,不是运行时反射所有 public 方法。常见失效点:
消费者类没加
特性,或特性值和发布时的 topic 名不完全一致(区分大小写)
方法签名必须是
或
,且参数只能是
、
、
或具体 DTO 类型(需匹配 JSON 序列化结构),多一个
参数就会跳过注册
项目未引用
(或 RabbitMQ/SQL Server 包),或
时没调用对应扩展方法如
,会导致
找不到任何订阅者
消息重复消费怎么防?CAP 本身不保证 Exactly-Once
CAP 的重试机制是“至少一次(At-Least-Once)”,网络抖动、消费者进程崩溃、超时未响应都会触发重发。业务层必须自己实现幂等:
C知道
CSDN推出的一款AI技术问答工具
下载
推荐在消费者方法开头查
表,用
做唯一索引去重(CAP 自动写该表,但需确保你没禁用
)
避免用“订单号”作为幂等键——如果发布方重发了两次不同内容但相同订单号的消息,后一次会覆盖前一次状态
不要依赖数据库主键冲突报错来判断重复,因为 CAP 可能因配置错误导致
表未启用,此时唯一防护就是业务逻辑内嵌判重(如更新时加
)
调度器卡住不动?看 Bootstrapper 的 ExecuteAsync 是否完成
CAP 启动本质是
实现了
,它的
方法必须成功返回,否则
不会开始轮询。典型阻塞原因:
数据库连接字符串指向了不可达实例,
无限重试(默认 5 秒间隔,无上限)
调用
时传入了空
,导致
或
未初始化,后续
获取不到发送器实例而静默失败
日志级别设为
以上,会掩盖
内部的
错误,实际需开
级别才能看到轮询频率和异常堆栈
最易被忽略的是:CAP 的“事务绑定”只对 EF Core 的
生效,如果你用 Dapper 或原生 ADO.NET 手动提交事务,
写入的记录不会和业务数据共用同一事务——必须显式传入
到
的
参数里。
PublishAsyncCap.PublishedProcessingDispatcherStorageOptions.ConnectionStringPublishAsyncInvalidOperationExceptionusing var tx = await _context.Database.BeginTransactionAsync()PublishAsyncCommitPublishAsyncheaderscap-msg-idcap-msg-name[CapSubscribe("xxx.order.created")]async TaskTaskCapMessagestringbyte[]ILoggerDotNetCore.CAP.KafkaAddCap.UseKafka(...)ConsumerRegisterCap.ReceivedMessageIdEnableReceivedLogCap.ReceivedWHERE status = 'pending'BootstrapperIHostedServiceExecuteAsyncDispatcherStorageInitializer.InitializeAsyncAddCapsetupActionKafkaOptionsSqlServerOptionsDispatcherWarningDispatcherFailed to publish messageDebugSaveChangesAsyncPublishAsyncDbContext.Database.CurrentTransactionPublishAsynctransaction