MassTransit中RoutingSlipActivityCompleted消费者未触发问题求助
问题:RoutingSlipActivityCompleted消费者未触发
我实现了包含Routing Slip构建逻辑的代码,同时编写了消费RoutingSlipActivityCompleted消息的RoutingSlipEventConsumer,但该消费者从未被触发。以下是相关代码:
Routing Slip构建代码
public class UpsertCredentialConsumer : IConsumer<IUpsertCredentialRequest>, IConsumer<RoutingSlipCompleted>, IConsumer<RoutingSlipFaulted> { private readonly ILogger<UpsertCredentialConsumer> _logger; public UpsertCredentialConsumer(ILogger<UpsertCredentialConsumer> logger) { _logger = logger; } public async Task Consume(ConsumeContext<IUpsertCredentialRequest> context) { var routingSlip = CreateRoutingSlip(context); await context.Execute(routingSlip); } private RoutingSlip CreateRoutingSlip(ConsumeContext<IUpsertCredentialRequest> context) { var builder = new RoutingSlipBuilder(NewId.NextGuid()); builder.SetVariables(new { context.RequestId, context.ResponseAddress }); foreach(var cred in context.Message.AgdtCredential) { //Se credential non presente if (cred.AsIsCredential is null) { //InsertCredential AddInsertCredentialActivity(builder, context.Message, cred.Credential); //InsertExternalRefeCredential builder.AddActivity(XXX.Name, XXX.Queue, new { ... }); } //Se credential presente if (cred.AsIsCredential is not null) { //UpdateCredential AddUpdateCredentialActivity(builder, context.Message, cred.Credential, cred.AsIsCredential); if (!cred.Credential.Status.Equals(cred.AsIsCredential.Status)) { //SetCredentialStatus se lo stato da settare è diverso da quello esistente builder.AddActivity(YYY.Name, YYY.Queue, new { ... }); } } else { //SetCredentialStatus builder.AddActivity(YYY.Name, YYY.Queue, new { ... }); } } builder.AddSubscription(context.ReceiveContext.InputAddress, RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted & RoutingSlipEvents.ActivityCompleted); return builder.Build(); } public async Task Consume(ConsumeContext<RoutingSlipCompleted> context) { _logger.LogTrace("(IUpsertCredentialRequest) Routing Slip Completed: {TrackingNumber}", context.Message.TrackingNumber); var requestId = context.GetVariable<Guid>("RequestId"); var responseAddress = context.GetVariable<Uri>("ResponseAddress"); var responseEndpoint = await context.GetSendEndpoint(responseAddress); await responseEndpoint.Send<IUpsertCredentialResult>( new { Success = true, __RequestId = requestId, Data = context.Message.Variables["credentialId"] }); } public async Task Consume(ConsumeContext<RoutingSlipFaulted> context) { _logger.LogTrace("(IUpsertCredentialRequest) Routing Slip Faulted: {TrackingNumber} {ExceptionInfo}", context.Message.TrackingNumber, context.Message.ActivityExceptions.FirstOrDefault()); 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 object errorValidation); var exceptions = context.Message.ActivityExceptions.Select(x => x.ExceptionInfo); var errorCode = errorValidation is null ? exceptions.Select(x => x.Message).FirstOrDefault() : ScsExceptions.ScsValidationError.ToString(); _logger.LogTrace("{TrackingNumber} ErrorCode: {error}", context.Message.TrackingNumber, errorCode); await responseEndpoint.Send<IUpsertCredentialResult>( new { Success = false, ErrorCode = errorCode, __RequestId = requestId }); } } }
RoutingSlipActivityCompleted消费者代码
public class RoutingSlipEventConsumer : IConsumer<RoutingSlipActivityCompleted>, IConsumer<RoutingSlipFaulted> { private readonly ILogger<RoutingSlipEventConsumer> _logger; public RoutingSlipEventConsumer(ILogger<RoutingSlipEventConsumer> logger) { _logger = logger; } public Task Consume(ConsumeContext<RoutingSlipActivityCompleted> context) { _logger.LogInformation("Routing Slip Activity Completed: {TrackingNumber} {ActivityName}", context.Message.TrackingNumber, context.Message.ActivityName); return Task.CompletedTask; } public Task Consume(ConsumeContext<RoutingSlipFaulted> context) { _logger.LogInformation("Routing Slip Faulted: {TrackingNumber} {ExceptionInfo}", context.Message.TrackingNumber, context.Message.ActivityExceptions.FirstOrDefault()); return Task.CompletedTask; } }
消费者注册代码
services.AddMassTransit(x => { x.AddConsumer<allTheOtherConsumer>().Endpoint(x=> x.Temporary = true); x.AddConsumer<RoutingSlipEventConsumer>().Endpoint(x => x.Temporary = true); x.UsingRabbitMq((context, cfg) => { cfg.Host(...); cfg.ConfigureEndpoints(context); }); });
排查原因及解决方案
1. 事件订阅的位运算错误
你使用了RoutingSlipEvents.Completed & RoutingSlipEvents.Faulted & RoutingSlipEvents.ActivityCompleted,这是位与运算,结果为0(三个事件是不同的位标志),相当于没有订阅任何事件。正确的应该使用**位或运算|**来合并多个事件:
builder.AddSubscription(目标地址, RoutingSlipEvents.Completed | RoutingSlipEvents.Faulted | RoutingSlipEvents.ActivityCompleted);
2. 订阅地址指向错误
当前代码中用context.ReceiveContext.InputAddress作为订阅地址,这个地址是UpsertCredentialConsumer的接收队列,而非RoutingSlipEventConsumer的队列,导致RoutingSlipActivityCompleted消息发送到了错误的队列。
解决办法:
获取RoutingSlipEventConsumer的队列地址进行订阅:
private async Task<RoutingSlip> CreateRoutingSlip(ConsumeContext<IUpsertCredentialRequest> context) { var builder = new RoutingSlipBuilder(NewId.NextGuid()); // ... 其他代码 // 获取RoutingSlipEventConsumer的端点地址 var consumerQueueName = KebabCaseEndpointNameFormatter.Instance.Consumer<RoutingSlipEventConsumer>(); var consumerEndpoint = await context.GetSendEndpoint(new Uri($"queue:{consumerQueueName}")); builder.AddSubscription(consumerEndpoint.Address, RoutingSlipEvents.ActivityCompleted | RoutingSlipEvents.Faulted | RoutingSlipEvents.Completed); return builder.Build(); }
3. 临时队列的生命周期问题
你给RoutingSlipEventConsumer配置了Temporary = true,临时队列是动态生成的,且应用停止后会被删除。如果订阅时没有正确获取到这个动态地址,或者队列未提前创建,消息无法送达。
解决办法:
移除临时队列配置,使用持久化队列:
x.AddConsumer<RoutingSlipEventConsumer>(); // 去掉.Endpoint(x => x.Temporary = true)
4. 辅助排查建议
- 检查RabbitMQ控制台,查看
RoutingSlipActivityCompleted消息是否被发送到正确的队列,是否有未被消费的消息。 - 调整日志级别为Debug/Trace,查看MassTransit的内部日志,确认消息的发送和路由情况。
内容的提问来源于stack exchange,提问作者Francesco Vargas
相关产品推荐
相关产品推荐

