Mass Transit Courier超时未消费时的分布式事务补偿问题
解决Mass Transit Courier活动超时未触发补偿的问题
要解决这个问题,核心是给你的Routing Slip活动配置超时机制,同时确保每个活动都有对应的补偿逻辑,并且Mass Transit能正确调度超时触发的补偿流程。下面是具体的实现步骤:
1. 确保RabbitMQ支持延迟消息
Mass Transit Courier的超时功能依赖于延迟消息调度,所以你需要先安装RabbitMQ的rabbitmq_delayed_message_exchange插件。安装完成后重启RabbitMQ生效。
2. 配置Bus启用延迟消息调度
在创建Bus实例时,启用延迟交换作为消息调度器,这样Courier才能在超时后触发补偿流程:
var bus = Bus.Factory.CreateUsingRabbitMq(cfg => { var host = cfg.Host(new Uri(RabbitMqConstants.RabbitMqUri), h => { h.Username("your-username"); h.Password("your-password"); }); // 启用延迟交换消息调度,支持超时触发 cfg.UseDelayedExchangeMessageScheduler(); // 注册执行活动和对应的补偿活动 cfg.ReceiveEndpoint("Core_Coding_Insert", e => { // 注册执行活动 e.ExecuteActivity<CoreCodingInsertActivity, Coding>(); // 注册补偿活动(必须实现ICompensateActivity接口) e.CompensateActivity<CoreCodingInsertCompensateActivity, Coding>(); }); cfg.ReceiveEndpoint("Kachar_Coding_Insert", e => { e.ExecuteActivity<KacharCodingInsertActivity, Coding>(); e.CompensateActivity<KacharCodingInsertCompensateActivity, Coding>(); }); // 其他活动同理注册... });
3. 给活动设置超时时间
你可以给单个活动设置超时,也可以给整个Routing Slip设置全局超时,两种方式可以结合使用:
单个活动超时
在添加活动时,通过ActivityOptions指定超时时长:
public async Task<Guid> PublishInsertCoding(Coding coding) { var builder = new RoutingSlipBuilder(NewId.NextGuid()); builder.AddActivity("Core_Coding_Insert", new Uri($"{RabbitMqConstants.RabbitMqUri}Core_Coding_Insert")); // 给Kachar的活动设置5分钟超时 builder.AddActivity("Kachar_Coding_Insert", new Uri($"{RabbitMqConstants.RabbitMqUri}Kachar_Coding_Insert"), options => options.Timeout = TimeSpan.FromMinutes(5)); builder.AddActivity("Rahavard_Coding_Insert", new Uri($"{RabbitMqConstants.RabbitMqUri}Rahavard_Coding_Insert")); builder.SetVariables(coding); var routingSlip = builder.Build(); await _bus.Execute(routingSlip); return routingSlip.TrackingNumber; }
全局Routing Slip超时
如果希望整个事务有一个总超时时间(不管单个活动状态),可以在构建Routing Slip时设置:
builder.SetTimeout(TimeSpan.FromMinutes(10)); // 整个事务10分钟超时
4. 实现补偿活动
每个执行活动都需要对应的补偿逻辑,当活动超时(或执行失败)时,Courier会自动触发补偿流程,对已经成功执行的活动按逆序进行补偿。
比如Kachar_Coding_Insert的补偿活动实现:
public class KacharCodingInsertCompensateActivity : ICompensateActivity<Coding> { private readonly IKacharRepository _repository; public KacharCodingInsertCompensateActivity(IKacharRepository repository) { _repository = repository; } public async Task Compensate(CompensateContext<Coding> context) { // 这里写具体的补偿逻辑:比如删除已经插入的Coding记录 await _repository.DeleteCodingAsync(context.Data.Id); // 记录补偿日志,方便后续排查问题 context.LogInformation($"补偿Kachar_Coding_Insert完成,ID: {context.Data.Id}"); } }
工作原理说明
当某个活动超过指定时间未被消费执行时,Mass Transit的调度器会触发Routing Slip的补偿流程:
- 自动识别哪些活动已经成功执行
- 按照与执行顺序相反的顺序调用对应的补偿活动
- 所有补偿完成后,整个分布式事务标记为已补偿
内容的提问来源于stack exchange,提问作者Meysam Gheysaryan
相关产品推荐
相关产品推荐

