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

MassTransit中RoutingSlipActivityCompleted入跳过队列致请求超时求助

问题:MassTransit+RabbitMQ路由单请求响应超时,RoutingSlipActivityCompleted进入跳过队列

在RabbitMQ上使用MassTransit,通过请求/响应模式启动包含多活动的路由单,要求路由单完成时通知调用方。目前调用方触发TimeoutException无法接收响应,且发现RoutingSlipActivityCompleted消息进入跳过队列。项目此前运行正常,已检查消息绑定未发现异常。补充:经排查存在偶发启动时错误,询问原因。


相关代码片段

public class UpsertExternalRequestConsumer : IConsumer<IUpsertExternalRequest>, IConsumer<RoutingSlipCompleted>, IConsumer<RoutingSlipFaulted>


public async Task Consume(ConsumeContext<IUpsertExternalRequest> context)
{
    _logger.LogTrace("IUpsertExternalRequest RECEIVED: {@message}", context.Message);

    var routingSlip = CreateRoutingSlip(context);

    await context.Execute(routingSlip);
}

private RoutingSlip CreateRoutingSlip(ConsumeContext<IUpsertExternalRequest> context)
{
    var builder = new RoutingSlipBuilder(NewId.NextGuid());

    builder.SetVariables(new
    {
        context.RequestId,
        context.ResponseAddress
    });

     
    //InsertIdentity
    builder.AddActivity(InsertIdentityActivity.Name, InsertIdentityActivity.Queue,
    new
    {
        IdentityFields = ...
    });

    //InsertExtendedIdentity
    builder.AddActivity(InsertExtendedIdentityActivity.Name, InsertExtendedIdentityActivity.Queue, 
new
    {
        IdentityFields = EntitiesExtensions.MapToKeyValues(context.Message.Identity.AdditionalInfo, BtpExternalFieldsMapping._ExtIdentitydict)
    });

    //other activities...
    

    builder.AddSubscription(context.ReceiveContext.InputAddress, RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted);
 
    return builder.Build();
}


public async Task Consume(ConsumeContext<RoutingSlipCompleted> context)
{
    var requestId = context.GetVariable<Guid>("RequestId");
    var responseAddress = context.GetVariable<Uri>("ResponseAddress");

    var responseEndpoint = await context.GetSendEndpoint(responseAddress);

    await responseEndpoint.Send<IUpsertExternalResponse>(
        new
        {
            Success = true,
            __RequestId = requestId,
            Data = context.Message.Variables["identityId"]
        });
}

public async Task Consume(ConsumeContext<RoutingSlipFaulted> context)
{
    var requestId = context.GetVariable<Guid>("RequestId");
    var responseAddress = context.GetVariable<Uri>("ResponseAddress");

    if (requestId.HasValue && responseAddress != null)
    {
        var responseEndpoint = await context.GetSendEndpoint(responseAddress);

        var validationError = context.Message.Variables.TryGetValue("ErrorCode", out string errorValidation);

        var exceptions = context.Message.ActivityExceptions.Select(x => x.ExceptionInfo);

        _logger.LogTrace("RoutingSlipFaulted: sending IBaseResponse. requestId:{requestId}, ResponseAddress:{responseAddress}", requestId, responseAddress);

        await responseEndpoint.Send<IUpsertExternalResponse>(
                new
                {
                    Success = false,
                    ErrorCode = errorValidation is null ? exceptions.Select(x => x.Message).FirstOrDefault() : ScsExceptions.ScsValidationError.ToString(),
                    __RequestId = requestId,
                });
    }
}


//调用请求的代码:
 var response = await _requestClient.GetResponse<IUpsertExternalResponse>(new
 {
      //业务数据...
 }, cancellationToken);

跳过队列中的示例消息

{
  "messageId": "58060000-569c-0050-162a-08dc1e7bd20e",
  "requestId": null,
  "correlationId": "58060000-569c-0050-30cb-08dc1e7bd20c",
  "conversationId": "fc2b0000-569c-0050-525f-08dc1e7bd1ec",
  "initiatorId": "fc2b0000-569c-0050-8abe-08dc1e7bd1f0",
  "sourceAddress": "rabbitmq://****/update-extended-identity_execute",
  "destinationAddress": "rabbitmq://****/UpsertExternalRequest",
  "responseAddress": null,
  "faultAddress": null,
  "messageType": [
    "urn:message:MassTransit.Courier.Contracts:RoutingSlipActivityCompleted"
  ],
  "message": {
    "trackingNumber": "fc2b0000-569c-0050-8abe-08dc1e7bd1f0",
    "timestamp": "2024-01-26T14:33:53.3592779Z",
    "duration": "00:00:00.0122791",
    "executionId": "58060000-569c-0050-30cb-08dc1e7bd20c",
    "activityName": "UpdateExtendedIdentity",
    "host": {
      "machineName": "IT4076YW",
      "processName": "****",
      "processId": 1624,
      "assembly": "****",
      "assemblyVersion": "1.0.9.0",
      "frameworkVersion": "6.0.26",
      "massTransitVersion": "8.0.13.0",
      "operatingSystemVersion": "Microsoft Windows NT 10.0.17763.0"
    },
    "arguments": {
      "loggedUser": 289,
      "appCode": "****",
      "identityId": 124288,
      "identityFields": [
        //业务数据... 
                
    ],
    "data": {
      "identityId": 124288,
      //*** 
    },
    "variables": {
      "requestId": "fc2b0000-569c-0050-41fc-08dc1e7bd1e9",
      "responseAddress": "rabbitmq://****Api_bus_9oioyynsuoyfy9z6bdqbh6ofrt?temporary=true",
      "identityId": 124288,
      "idLoggedUser": 289
    }
  },
  "expirationTime": null,
  "sentTime": "2024-01-26T14:33:53.3717034Z",
  "headers": {},
  "host": {
    "machineName": "IT4076YW",
    "processName": "****",
    "processId": 1624,
    "assembly": "****",
    "assemblyVersion": "1.0.9.0",
    "frameworkVersion": "6.0.26",
    "massTransitVersion": "8.0.13.0",
    "operatingSystemVersion": "Microsoft Windows NT 10.0.17763.0"
  }
}

补充问题

经排查发现偶发启动时错误,请问为何会出现这种偶发现象?


问题原因与解决方案

1. RoutingSlipActivityCompleted进入跳过队列的直接原因

你的UpsertExternalRequestConsumer订阅了RoutingSlipCompleted和RoutingSlipFaulted,但未订阅RoutingSlipActivityCompleted。当路由单单个活动完成时,MassTransit会发送该消息到订阅地址(UpsertExternalRequest队列),消费者无法处理,因此被转入跳过队列。

2. 请求超时(TimeoutException)的核心问题

  • 订阅配置错误:builder.AddSubscription(...)中使用了RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted位运算,这会导致结果为0,等于未订阅任何事件,RoutingSlipCompleted消息根本不会发送到你的消费者,这是超时的核心原因。应改为RoutingSlipEvents.Completed | RoutingSlipEvents.Faulted(逻辑或)。
  • 临时响应地址过期:请求客户端创建的临时响应地址带有?temporary=true,存在生命周期限制。如果路由单执行时间超过临时队列过期时间,队列被自动删除,响应无法送达。
  • 消费方法未捕获异常:Consume(ConsumeContext<RoutingSlipCompleted>)中获取identityId时,若变量不存在会抛出未捕获异常,导致消息重试或进入死信队列,无法发送响应。

3. 偶发启动错误的可能原因

  • RabbitMQ连接初始化延迟:启动时应用未完全建立RabbitMQ连接,路由单活动的队列绑定未完成,消息无法被消费,引发后续超时。
  • 拓扑结构不一致:多实例部署时,不同实例的MassTransit拓扑存在差异,启动时竞争资源导致部分绑定失败。
  • 依赖服务未就绪:路由单依赖的外部服务(如数据库、身份服务)在启动时未就绪,导致活动执行失败触发异常。

修复步骤

  1. 修正订阅事件运算:将RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted改为RoutingSlipEvents.Completed | RoutingSlipEvents.Faulted,确保订阅完成和故障事件。
  2. 处理RoutingSlipActivityCompleted消息:要么在消费者中添加IConsumer<RoutingSlipActivityCompleted>实现(可空实现),要么在订阅时排除该事件,避免消息进入跳过队列。
  3. 延长临时队列生命周期:创建RequestClient时,配置更长的RequestTimeout和TemporaryQueueExpiration,避免临时队列提前过期。
  4. 增加异常捕获:在RoutingSlipCompleted和RoutingSlipFaulted消费方法中添加try-catch块,记录异常并确保响应能被发送。
  5. 启动健康检查:应用启动时添加RabbitMQ连接、依赖服务的健康检查,确保所有服务就绪后再处理请求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 04:55:59