You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的补偿流程:

  1. 自动识别哪些活动已经成功执行
  2. 按照与执行顺序相反的顺序调用对应的补偿活动
  3. 所有补偿完成后,整个分布式事务标记为已补偿

内容的提问来源于stack exchange,提问作者Meysam Gheysaryan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.07 12:22:43