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

如何在MassTransit中实现SagaConsumeContext及过滤器注册?

监听MassTransit Saga处理事件的过滤器实现与注册

我需要监听Saga(SagaStateMachineInstance)处理的事件消息,目前已实现通过ConsumerConsumeContext监听MassTransit普通消费者处理消息的方案,但对于Saga,无法理解对应的过滤器及在MassTransit中的注册方式,已编写如下过滤器代码:

public class SagaConsumeFilter<TSaga, TMessage> : IFilter<SagaConsumeContext<TSaga, TMessage>> 
    where TSaga : class, ISaga 
    where TMessage : class
{
    public async Task Send(SagaConsumeContext<TSaga, TMessage> context, IPipe<SagaConsumeContext<TSaga, TMessage>> next)
    {
        Console.WriteLine($"[SAGA FILTER] Saga: {typeof(TSaga).Name}");
        Console.WriteLine($"[SAGA FILTER] Message: {typeof(TMessage).Name}");
        Console.WriteLine($"[SAGA FILTER] CorrelationId: {context.Saga.CorrelationId}");

        await next.Send(context);
    }

    public void Probe(ProbeContext context)
    {
        context.CreateFilterScope("SagaConsumeFilter");
    }
}

过滤器注册方案

MassTransit中Saga的过滤器注册逻辑和普通消费者不同,需通过Saga专属配置管道添加,以下是三种可行的注册方式:

1. 针对单个Saga注册过滤器

在注册指定Saga时,通过ConfigureSaga方法直接挂载过滤器:

services.AddMassTransit(x =>
{
    x.AddSaga<YourSagaClass>()
      .InMemoryRepository(); // 可替换为EF Core/MongoDB等实际存储实现

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.ConfigureEndpoints(context);

        // 为目标Saga配置过滤器
        cfg.ConfigureSaga<YourSagaClass>(context)
           .UseFilter(new SagaConsumeFilter<YourSagaClass, object>()); // 使用object适配该Saga处理的所有消息类型
    });
});

2. 全局注册所有Saga的过滤器

若需对系统内所有Saga统一应用过滤器,可创建通用过滤器定义类并全局注册:

// 全局过滤器定义类
public class SagaConsumeFilterDefinition : SagaFilterDefinition
{
    protected override void ConfigureFilter(IFilterConfigurator<SagaConsumeContext> configurator)
    {
        configurator.UseFilter(new SagaConsumeFilter<ISaga, object>());
    }
}

然后在MassTransit配置中添加全局过滤器:

services.AddMassTransit(x =>
{
    x.AddSaga<YourSagaClass>()
      .InMemoryRepository();

    // 注册全局Saga过滤器
    x.AddSagaFilter<SagaConsumeFilterDefinition>();

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.ConfigureEndpoints(context);
    });
});

3. 针对特定消息类型的Saga事件注册过滤器

如果仅需监听Saga处理某类特定消息的流程,可在状态机的事件配置中单独指定过滤器:

public class YourSagaStateMachine : MassTransitStateMachine<YourSagaClass>
{
    public YourSagaStateMachine()
    {
        // 状态定义...

        Event(() => OrderSubmittedEvent, e =>
        {
            e.Filter(new SagaConsumeFilter<YourSagaClass, OrderSubmitted>());
            // 其他事件处理配置...
        });
    }

    public Event<OrderSubmitted> OrderSubmittedEvent { get; set; }
}

注意事项

  • 泛型过滤器的类型参数需匹配Saga和消息类型,避免运行时类型转换异常;
  • 过滤器中的Probe方法用于MassTransit监控探针,建议保留以支持系统监控功能;
  • 不同Saga存储提供者(如InMemory、EF Core)的注册细节略有差异,但过滤器添加逻辑一致。

内容的提问来源于stack exchange,提问作者Ola Hällvall

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 18:43:10