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

MassTransit发布过滤器能否替换待发送的原始消息?

问题

我们使用MassTransit对接Azure Service Bus,但Service Bus存在最大消息大小限制,部分消息超出了该限制。希望实现一个发布过滤器完成以下操作:

  • 将原始消息存入存储服务
  • 用同类型的“空”实例替换原始消息
  • 发送替换后的消息而非原始消息

消费端收到消息后,会通过消息ID从存储服务获取完整消息体。

尝试使用PublishContextProxy替换消息,但未生效——发布的仍是原始消息,示例代码如下:

public class MyEventEventPublishFilter<T> : IFilter<PublishContext<T>>
where T : class 
{
    private readonly IStorageService _storage;

    public CampaignTraffickingEventPublishFilter(IStorageService storage)
    {
        _storage = storage;
    }

    public async Task Send(PublishContext<T> context, IPipe<PublishContext<T>> next)
    {
        if (context is PublishContext<MyBigEvent> { Message: var msg })
        {
            await _storageService.SetAsync(context.MessageId, evt);

            var newContext = new PublishContextProxy<MyBigEvent>(
                context,
                new MyBigEvent(evt.Id, null, null, 0, null, null, null, null)); // 尝试替换消息
            await next.Send((newContext as PublishContextProxy<T>)!);
            evt.Fields = null;

            return;
        }

        await next.Send(context);
    }
}

结果原始的“非空”消息仍被发布到主题中。

解决方案

MassTransit完全支持发送前替换消息的场景,你代码的问题在于PublishContextProxy的使用逻辑错误——它没有正确覆盖原始上下文的消息序列化流程。

正确的做法是直接修改当前上下文的Message属性,因为PublishContext<T>中的Message是可写的,替换后让管道继续处理修改后的上下文即可。

调整后的代码示例:

public class MyEventPublishFilter<T> : IFilter<PublishContext<T>>
where T : class 
{
    private readonly IStorageService _storageService;

    public MyEventPublishFilter(IStorageService storageService)
    {
        _storageService = storageService;
    }

    public async Task Send(PublishContext<T> context, IPipe<PublishContext<T>> next)
    {
        if (context.Message is MyBigEvent evt)
        {
            // 1. 将原始消息存入存储服务
            await _storageService.SetAsync(context.MessageId, evt);
            
            // 2. 创建同类型的占位实例
            var placeholder = new MyBigEvent(evt.Id, null, null, 0, null, null, null, null);
            
            // 3. 直接替换上下文的Message属性
            context.Message = placeholder as T;
            
            // 继续执行管道,发送替换后的消息
            await next.Send(context);
            return;
        }

        await next.Send(context);
    }
}

额外注意事项:

  • 确保MyBigEvent的序列化逻辑能正确处理占位实例,避免序列化后仍超出大小限制
  • 消费端需做好异常处理,比如存储服务不可用、消息ID不存在等场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 20:59:56