如何根据条件创建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
相关产品推荐
相关产品推荐

