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

MassTransit/RabbitMQ:如何获取跳过队列消息并记录日志?

解决MassTransit中消息进入Skipped队列并通过Serilog记录的测试方案

一、主动触发消息进入Skipped队列(最直接的测试方式)

在现有业务消费者中,通过调用context.Skip()方法主动跳过符合特定条件的消息,MassTransit会自动将该消息转发到全局的masstransit-skipped队列。

示例业务消费者代码

public class MyMessageConsumer : IConsumer<MyMessage>
{
    private readonly ILogger<MyMessageConsumer> _logger;

    public MyMessageConsumer(ILogger<MyMessageConsumer> logger)
    {
        _logger = logger;
    }

    public async Task Consume(ConsumeContext<MyMessage> context)
    {
        // 自定义跳过条件:比如消息ID匹配测试标识
        if (context.Message.Id.Equals("test-skip-flag", StringComparison.OrdinalIgnoreCase))
        {
            _logger.LogInformation("Initiating skip for message {MessageId}", context.MessageId);
            await context.Skip("Test skip trigger"); // 可传入自定义跳过原因
            return;
        }

        // 正常业务处理逻辑
        _logger.LogInformation("Processing message {MessageId}", context.MessageId);
        // ... 你的业务代码
    }
}

二、注册Skipped队列的消费者

创建专门的消费者监听masstransit-skipped队列,接收跳过的消息并通过Serilog记录详情:

public class SkippedMessageConsumer : IConsumer<SkippedMessage>
{
    private readonly ILogger<SkippedMessageConsumer> _logger;

    public SkippedMessageConsumer(ILogger<SkippedMessageConsumer> logger)
    {
        _logger = logger;
    }

    public async Task Consume(ConsumeContext<SkippedMessage> context)
    {
        var skippedMsg = context.Message;
        _logger.Warning(
            "Skipped message tracked - Original Message ID: {MsgId}, Reason: {Reason}, Timestamp: {Time}",
            skippedMsg.MessageId,
            skippedMsg.Reason,
            skippedMsg.Timestamp);

        await Task.CompletedTask;
    }
}

三、MassTransit与Serilog的配置(.NET 6)

在Program.cs中完成服务注册,确保Serilog作为日志提供者,同时注册两个消费者:

// 配置Serilog
builder.Host.UseSerilog((ctx, config) =>
{
    config.ReadFrom.Configuration(ctx.Configuration)
          .Enrich.FromLogContext()
          .WriteTo.Console() // 控制台输出
          .WriteTo.File("logs/skipped-messages.log", rollingInterval: RollingInterval.Day); // 按天生成日志文件
});

// 注册MassTransit
builder.Services.AddMassTransit(x =>
{
    // 注册业务消费者和跳过消息消费者
    x.AddConsumer<MyMessageConsumer>();
    x.AddConsumer<SkippedMessageConsumer>();

    x.UsingRabbitMq((ctx, cfg) =>
    {
        cfg.Host("rabbitmq://localhost", h =>
        {
            h.Username("guest");
            h.Password("guest");
        });

        // 自动配置所有消费者的端点
        cfg.ConfigureEndpoints(ctx);

        // 可选:配置全局异常跳过策略(特定异常触发自动跳过)
        cfg.ReceiveEndpoint("my-business-queue", e =>
        {
            e.Consumer<MyMessageConsumer>(ctx);
            // 重试2次后,若抛出InvalidOperationException则自动跳过
            e.UseMessageRetry(r => r.Interval(2, 1000));
            e.SkipMessageOnException<InvalidOperationException>();
        });
    });
});

四、测试发送跳过消息

通过发布者发送符合跳过条件的消息,即可触发完整流程:

public class TestMessagePublisher
{
    private readonly IPublishEndpoint _publishEndpoint;

    public TestMessagePublisher(IPublishEndpoint publishEndpoint)
    {
        _publishEndpoint = publishEndpoint;
    }

    public async Task SendSkipTestMessage()
    {
        await _publishEndpoint.Publish(new MyMessage
        {
            Id = "test-skip-flag",
            Content = "This message will be skipped"
        });
    }
}

额外:直接发送消息到Skipped队列(仅测试用)

如果需要跳过业务流程直接测试Skipped队列的消费逻辑,可以直接发送SkippedMessage到对应的交换器:

public async Task DirectSendToSkippedQueue(ISendEndpointProvider sendEndpointProvider)
{
    var skippedEndpoint = await sendEndpointProvider.GetSendEndpoint(new Uri("exchange:masstransit-skipped"));
    await skippedEndpoint.Send(new SkippedMessage
    {
        MessageId = Guid.NewGuid().ToString(),
        Reason = "Direct test skip",
        Timestamp = DateTime.UtcNow
    });
}

内容的提问来源于stack exchange,提问作者Tom van den Bogaart

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 15:40:24