如何修改MassTransit传输至_error队列的消息内容?
实现方案
MassTransit默认的故障投递逻辑会将原始消费消息完整发送到错误队列,要替换为自定义的错误消息,可参考以下两种实现方案:
方案1:消费逻辑内直接发送自定义错误消息(适配当前场景,推荐)
不需要抛出异常触发默认错误投递流程,构造好自定义错误消息后手动发送到错误队列即可,修改后的消费者代码如下:
public class MessageConsumer : IConsumer<Message> { public MessageConsumer() {} public async Task Consume(ConsumeContext<Message> context) { if (context.Message.Records.Any(m => m.Contains("Fault"))) { var faultedRecords = context.Message.Records.Where(r => r.Contains("Fault")).ToList(); // 仅包含异常记录的消息 var errorMessage = new Message() { Records = faultedRecords }; // 手动发送到错误队列,_error队列地址规则为 接收端点地址+"_error",也可提前配置为固定地址 var errorQueueAddress = new Uri($"{context.ReceiveContext.InputAddress}_error"); await context.Send(errorQueueAddress, errorMessage); // 直接返回标记消费成功,无需抛异常触发默认错误流程 return; } //...其他业务逻辑 return Task.CompletedTask; } }
该方案逻辑直观,无需修改全局管线配置,仅调整当前消费者逻辑即可生效。
方案2:全局异常过滤器统一处理(适合多消费者通用规则场景)
如果需要全局生效,所有消费者抛出异常时统一替换错误消息内容,可实现IFilter<ReceiveContext>自定义异常处理逻辑,在消息进入错误队列前替换Payload:
- 首先实现自定义过滤器
public class CustomErrorFilter : IFilter<ReceiveContext> { public async Task Send(ReceiveContext context, IPipe<ReceiveContext> next) { try { await next.Send(context); } catch (Exception ex) { // 判断当前消费的是Message类型消息 if (context.TryGetMessage<Message>(out var messageContext)) { var originalMessage = messageContext.Message; // 构造自定义错误消息 var faultedRecords = originalMessage.Records.Where(r => r.Contains("Fault")).ToList(); var errorMessage = new Message() { Records = faultedRecords }; // 替换上下文的消息体 context.SetMessage(errorMessage); } // 继续抛出异常,触发MassTransit默认的错误队列投递逻辑 throw; } } public void Probe(ProbeContext context) { context.CreateFilterScope("custom-error-filter"); } }
- 在接收端点配置中注册过滤器
cfg.ReceiveEndpoint("你的业务队列名称", e => { e.Consumer<MessageConsumer>(); // 注册自定义错误过滤器 e.UseFilter(new CustomErrorFilter()); });
内容的提问来源于stack exchange,提问作者Andrei Kleban
相关产品推荐
相关产品推荐

