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
相关产品推荐
相关产品推荐

