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

RoutingSlip无法回传至发布者的超时问题求助

问题:MassTransit RoutingSlip完成后回传结果至发布者超时

场景复现

发布者通过GetResponse发起请求获取支付处理结果;消费者接收请求后创建RoutingSlip执行一系列活动;RoutingSlip完成后,在RoutingSlipCompleted消费者中尝试回传结果至发布者,触发如下超时错误:

{"Timeout waiting for response, RequestId: cc200000-fb75-4ced-c30e-08dbd87c612b"}

发布者代码

public async Task<IActionResult> AliEshghiPaymentUrl(GeneralAuthenticatePaymentRequestModel model)
{
    try
    {
        var result = await _massClientRequest.GetResponse<ResultPayMentInDBResponse>(new GeneralAuthenticatePaymentRequestModel
                {
                    MyRequestedUrl = model.MyRequestedUrl,
                    AccountRequest = model.AccountRequest,
                    UserId = model.UserId
                });

        return Ok(result);
    }
    catch (Exception ex)
    {
        throw;
    }
}

RoutingSlip创建消费者代码

public async Task Consume(ConsumeContext<GeneralAuthenticatePaymentRequestModel> context)
{
    var routingSlip = CreateRoutingSlip(context);
    await context.Execute(routingSlip);
}

private RoutingSlip CreateRoutingSlip(ConsumeContext<GeneralAuthenticatePaymentRequestModel> context)
{
    var builder = new RoutingSlipBuilder(NewId.NextGuid());
    builder.AddSubscription(context.ReceiveContext.InputAddress,
                            RoutingSlipEvents.Completed
                               | RoutingSlipEvents.Faulted);

    builder.AddVariable("userId", context.Message.UserId);
    builder.AddVariable(nameof(context.RequestId), context.RequestId);
    builder.AddVariable(nameof(context.ResponseAddress), context.ResponseAddress);
    builder.AddVariable(nameof(context.FaultAddress), context.FaultAddress);
    builder.AddVariable("Request", context.Message);

    var dbQueueName = _formatter.
                ExecuteActivity<CheckAuthenticatedRequestFromDBActivity, AuthenticatePaymentRequestModel>();
    builder.AddActivity(nameof(SetAuthenticatedRequestWithDbActivity),
                new Uri($"queue:{dbQueueName}"), 
                new {});
    var apiQueueName = _formatter.ExecuteActivity<DesideWhichBankToChoseActivity, PaymentBankRequestModel>();
    builder.AddActivity(nameof(DesideWhichBankToChoseActivity), new Uri($"queue:{apiQueueName}"), new
            {
                MyRequestUrl ="/myeshghi"
            });

    var addToDbQueueName = _formatter.ExecuteActivity<SetChoosedBankInDBActivity, SetPaymnetResultInDBRequest>();
    builder.AddActivity(nameof(SetChoosedBankInDBActivity), new Uri($"queue:{addToDbQueueName}"), new
            {
                MyRequestUrl = "/myeshghi"
            });
            
    return builder.Build();
}

RoutingSlipCompleted消费者代码

public async Task Consume(ConsumeContext<RoutingSlipCompleted> context)
{
    var isSucceed = context.GetVariable<bool>("IsSucceed");
    var model = new ResultPayMentInDBResponse();
    model.IsSucceed = isSucceed ==null? false: (bool)isSucceed;

    var requestId = context.GetVariable<Guid>(nameof(ConsumeContext.RequestId));
    var responseAddress = context.GetVariable<Uri>(nameof(ConsumeContext.ResponseAddress));

    if (requestId.HasValue && responseAddress != null)
    {
        var responseEndpoint = await context.GetSendEndpoint(responseAddress);
        await responseEndpoint.Send<ResultPayMentInDBResponse>(model);
    }
}

配置代码(Program.cs)

var builder = WebApplication.CreateBuilder(args);
builder.Services.AddMassTransit(x =>
{
    x.SetSnakeCaseEndpointNameFormatter();
    x.SetKebabCaseEndpointNameFormatter();
    x.AddConsumer<PaymentRequestFlowConsumer>();

    x.AddActivitiesFromNamespaceContaining<CourierActivities>();

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host(builder.Configuration["RabbitMQModel:HostAddress"], port:
            ushort.Parse(builder.Configuration["RabbitMQModel:Port"]), "/", h =>
        {
            h.Username(builder.Configuration["RabbitMQModel:Username"]);
            h.Password(builder.Configuration["RabbitMQModel:Password"]);
        });

        cfg.ReceiveEndpoint("PAYMENTFLOW_Request", e =>
        {
            e.ConfigureConsumer<PaymentRequestFlowConsumer>(context);
        });

        cfg.UseDelayedMessageScheduler();

        cfg.ConfigureEndpoints(context);
    });
});
var app = builder.Build();

app.MapGet("/", () => "business flow micro started!");

app.Run();

原因排查

  1. 默认超时时间过短:GetResponse默认超时为30秒,若RoutingSlip包含的数据库操作、第三方接口调用等活动总执行时间超过该值,发布者会提前终止等待,临时响应队列被销毁,后续回传无法找到目标队列。
  2. 回传未关联原始RequestId:发送响应时未指定原始请求的RequestId,发布者无法将响应与发起的请求匹配,导致持续等待。
  3. RoutingSlip活动执行阻塞/异常:单个活动执行卡住(如数据库死锁、第三方接口超时)或抛出未捕获异常,导致RoutingSlipCompleted事件触发延迟,超过发布者超时时间。
  4. 临时响应队列生命周期问题:发布者的临时响应队列在请求超时后自动销毁,即使后续RoutingSlipCompleted触发,也无法将消息发送到已销毁的队列。

解决办法

1. 延长GetResponse超时时间

调用GetResponse时显式指定更长的超时时间,匹配RoutingSlip的预估执行时长:

var result = await _massClientRequest.GetResponse<ResultPayMentInDBResponse>(
    new GeneralAuthenticatePaymentRequestModel
    {
        MyRequestedUrl = model.MyRequestedUrl,
        AccountRequest = model.AccountRequest,
        UserId = model.UserId
    }, 
    TimeSpan.FromMinutes(5)); // 根据实际业务调整时长

2. 回传时关联原始RequestId

发送响应时设置RequestId为原始请求的ID,确保发布者能匹配到对应的请求:

if (requestId.HasValue && responseAddress != null)
{
    var responseEndpoint = await context.GetSendEndpoint(responseAddress);
    await responseEndpoint.Send<ResultPayMentInDBResponse>(model, x =>
    {
        x.RequestId = requestId; // 关联原始请求ID
    });
}

3. 优化RoutingSlip活动执行效率

  • 排查并优化数据库查询、第三方接口调用的性能;
  • 为耗时活动添加超时控制,避免单个活动阻塞整个流程;
  • 检查活动代码,确保所有异常都被捕获并正确处理,避免RoutingSlip无法触发Completed事件。

4. 增加日志排查

在RoutingSlipCompleted消费者中添加日志,确认requestId、responseAddress的有效性,以及发送响应时是否有异常:

public async Task Consume(ConsumeContext<RoutingSlipCompleted> context)
{
    var logger = context.GetLogger<RoutingSlipCompletedConsumer>();
    try
    {
        var isSucceed = context.GetVariable<bool>("IsSucceed");
        var model = new ResultPayMentInDBResponse();
        model.IsSucceed = isSucceed ?? false;

        var requestId = context.GetVariable<Guid>(nameof(ConsumeContext.RequestId));
        var responseAddress = context.GetVariable<Uri>(nameof(ConsumeContext.ResponseAddress));

        logger.LogInformation("RoutingSlip completed, RequestId: {RequestId}, ResponseAddress: {ResponseAddress}", requestId, responseAddress);

        if (requestId.HasValue && responseAddress != null)
        {
            var responseEndpoint = await context.GetSendEndpoint(responseAddress);
            await responseEndpoint.Send<ResultPayMentInDBResponse>(model, x => x.RequestId = requestId);
            logger.LogInformation("Response sent successfully for RequestId: {RequestId}", requestId);
        }
        else
        {
            logger.LogWarning("Invalid RequestId or ResponseAddress for RoutingSlip: {SlipId}", context.Message.TrackingNumber);
        }
    }
    catch (Exception ex)
    {
        logger.LogError(ex, "Failed to send response for RoutingSlip: {SlipId}", context.Message.TrackingNumber);
        throw;
    }
}

5. 验证RoutingSlip订阅配置

确认RoutingSlipCompleted消费者已正确注册并能及时接收事件:

  • 在AddMassTransit中添加消费者注册:
x.AddConsumer<RoutingSlipCompletedConsumer>();
  • 确保ConfigureEndpoints自动创建对应的队列,或手动配置接收端点:
cfg.ReceiveEndpoint("routing-slip-completed", e =>
{
    e.ConfigureConsumer<RoutingSlipCompletedConsumer>(context);
});

内容的提问来源于stack exchange,提问作者Ali Eshghi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 15:19:50