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拓扑存在差异,启动时竞争资源导致部分绑定失败。
- 依赖服务未就绪:路由单依赖的外部服务(如数据库、身份服务)在启动时未就绪,导致活动执行失败触发异常。
修复步骤
- 修正订阅事件运算:将
RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted改为RoutingSlipEvents.Completed | RoutingSlipEvents.Faulted,确保订阅完成和故障事件。 - 处理RoutingSlipActivityCompleted消息:要么在消费者中添加
IConsumer<RoutingSlipActivityCompleted>实现(可空实现),要么在订阅时排除该事件,避免消息进入跳过队列。 - 延长临时队列生命周期:创建
RequestClient时,配置更长的RequestTimeout和TemporaryQueueExpiration,避免临时队列提前过期。 - 增加异常捕获:在
RoutingSlipCompleted和RoutingSlipFaulted消费方法中添加try-catch块,记录异常并确保响应能被发送。 - 启动健康检查:应用启动时添加RabbitMQ连接、依赖服务的健康检查,确保所有服务就绪后再处理请求。
内容的提问来源于stack exchange,提问作者dvRoss
相关产品推荐
相关产品推荐

