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

如何根据条件创建MassTransit Saga实例?

解决MassTransit Saga条件不满足仍创建空实例的问题

你的问题出在MassTransit处理Initially阶段事件的机制:当事件进入Initially分支时,框架会先自动创建一个新的Saga实例,再执行后续的IfAsync判断。哪怕条件不满足,这个已经创建的空实例还是会被持久化到数据库。

正确的实现方式

要避免这种情况,需要在Saga实例创建之前就过滤掉不符合条件的事件,有两种常用方案:


方案1:使用事件的FilterAsync过滤(推荐)

在事件配置阶段添加条件判断,不符合条件的事件直接被跳过,不会触发Saga实例的创建:

public MySaga()
{
    InstanceState(x => x.CurrentState, Processing);

    Event(() => EventOne, x =>
    {
        x.CorrelateById(context => context.Message.TransactionId);
        // 先判断条件,不满足则不处理该事件
        x.FilterAsync(async context => await ShouldCreateSaga(context));
    });

    Event(() => EventTwo, x =>
    {
        x.CorrelateById(context => context.Message.TransactionId);
        x.FilterAsync(async context => await ShouldCreateSaga(context));
    });

    Initially(
        When(EventOne)
            .Then(context =>
            {
                context.Saga.TransactionId = context.Message.TransactionId;
            })
            .TransitionTo(Processing),
        When(EventTwo)
            .Then(context =>
            {
                context.Saga.TransactionId = context.Message.TransactionId;
            })
            .TransitionTo(Processing)
    );

    // 其他状态逻辑...
}

private static async Task<bool> ShouldCreateSaga(PipeContext context)
{
    var serviceProvider = context.GetPayload<IServiceProvider>();
    var dbContext = serviceProvider.GetRequiredService<DbContext>();
    return await dbContext.SomeTable.AnyAsync(x => x.CreateSaga);
}

方案2:自定义OnMissingInstance逻辑

如果需要更灵活的实例创建控制,可以在事件配置中重写OnMissingInstance方法,手动控制是否创建实例:

public MySaga()
{
    InstanceState(x => x.CurrentState, Processing);

    Event(() => EventOne, x =>
    {
        x.CorrelateById(context => context.Message.TransactionId);
        x.OnMissingInstance(async context =>
        {
            // 先判断是否需要创建实例
            if (!await ShouldCreateSaga(context))
                return; // 不满足条件,直接忽略

            // 手动创建并初始化Saga实例
            var saga = new MySaga
            {
                CorrelationId = context.Message.TransactionId,
                TransactionId = context.Message.TransactionId,
                CurrentState = Processing
            };

            // 将实例添加到Saga仓库
            await context.SagaRepository.Add(context, saga);

            // 如需执行后续逻辑(如发布事件),可在此添加
        });
    });

    Event(() => EventTwo, x =>
    {
        x.CorrelateById(context => context.Message.TransactionId);
        x.OnMissingInstance(async context =>
        {
            if (!await ShouldCreateSaga(context))
                return;

            var saga = new MySaga
            {
                CorrelationId = context.Message.TransactionId,
                TransactionId = context.Message.TransactionId,
                CurrentState = Processing
            };

            await context.SagaRepository.Add(context, saga);
        });
    });

    // 其他状态逻辑...
}

private static async Task<bool> ShouldCreateSaga(PipeContext context)
{
    var serviceProvider = context.GetPayload<IServiceProvider>();
    var dbContext = serviceProvider.GetRequiredService<DbContext>();
    return await dbContext.SomeTable.AnyAsync(x => x.CreateSaga);
}

核心原理

两种方案都是把条件判断提前到Saga实例创建之前:

  • 方案1的FilterAsync会在事件进入Saga处理流程前直接过滤掉不符合条件的请求,完全不会触发实例创建逻辑。
  • 方案2通过OnMissingInstance接管了“找不到实例时”的处理逻辑,手动控制是否创建并初始化实例。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 01:04:55