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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 10:28:00