如何为IMessageScheduler添加适配IPublishEndpoint的发布过滤器?
让IntegrationEventPublishFilter支持IMessageScheduler定时消息的解决方案
问题描述
我编写了一个IntegrationEventPublishFilter发布过滤器,通过IPublishEndpoint发布非定时消息时可正常运行。现在需要让该过滤器在通过IMessageScheduler发布定时消息时也生效,现有代码如下:
public class IntegrationEventPublishFilter<T> : IFilter<PublishContext<T>> where T : class { public IntegrationEventPublishFilter() {} public Task Send(PublishContext<T> context, IPipe<PublishContext<T>> next) { // do something return next.Send(context); } public void Probe(ProbeContext context) { // do something } }
解决方案
过滤器对定时消息不生效的核心原因是:IMessageScheduler使用ScheduleContext<T>上下文(继承自PublishContext<T>),但现有过滤器仅实现了IFilter<PublishContext<T>>接口,且未注册到定时调度管道。需按以下步骤修改:
1. 更新过滤器实现,兼容ScheduleContext
让过滤器同时实现IFilter<ScheduleContext<T>>接口,复用核心过滤逻辑:
public class IntegrationEventPublishFilter<T> : IFilter<PublishContext<T>>, IFilter<ScheduleContext<T>> where T : class { public IntegrationEventPublishFilter() {} public Task Send(PublishContext<T> context, IPipe<PublishContext<T>> next) { // 核心过滤逻辑:比如添加消息头、记录事件日志等 return next.Send(context); } public Task Send(ScheduleContext<T> context, IPipe<ScheduleContext<T>> next) { // 复用PublishContext的处理逻辑,ScheduleContext是其直接子类 return Send((PublishContext<T>)context, next as IPipe<PublishContext<T>>); } public void Probe(ProbeContext context) { // 探针逻辑(可选,用于系统监控) var scope = context.CreateFilterScope("integration-event-publish"); scope.Add("message-type", typeof(T).Name); } }
2. 将过滤器注册到发布和调度管道
在MassTransit配置中,需要把过滤器同时添加到普通发布管道和定时调度管道:
方式一:开放泛型注册(推荐)
无需为每个消息类型单独配置,自动适配所有消息:
services.AddMassTransit(cfg => { // 其他配置:如定义消费者、设置消息传输等... // 注册到普通发布管道 cfg.AddPublishMessageFilter(typeof(IntegrationEventPublishFilter<>)); // 注册到定时调度管道 cfg.AddScheduleMessageFilter(typeof(IntegrationEventPublishFilter<>)); });
方式二:特定消息类型注册
仅针对指定消息类型启用过滤器:
services.AddMassTransit(cfg => { // 其他配置... cfg.UsingRabbitMq((ctx, cfg) => { // 配置普通发布管道 cfg.ConfigurePublish(c => { c.AddFilter(new IntegrationEventPublishFilter<YourTargetEvent>()); }); // 配置定时调度管道 cfg.ConfigureScheduling(c => { c.AddFilter(new IntegrationEventPublishFilter<YourTargetEvent>()); }); }); });
关键说明
ScheduleContext<T>继承自PublishContext<T>,因此可以直接复用已有过滤逻辑,无需重复编写- 必须同时注册到两个管道,才能让过滤器在普通消息发布和定时消息调度场景下均生效
内容的提问来源于stack exchange,提问作者KoKo
相关产品推荐
相关产品推荐

