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

