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

如何修改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:

  1. 首先实现自定义过滤器
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");
    }
}
  1. 在接收端点配置中注册过滤器
cfg.ReceiveEndpoint("你的业务队列名称", e =>
{
    e.Consumer<MessageConsumer>();
    // 注册自定义错误过滤器
    e.UseFilter(new CustomErrorFilter());
});

内容的提问来源于stack exchange,提问作者Andrei Kleban

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 15:57:02