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

