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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 09:25:54