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

MassTransit Saga OnMissingInstance配置异常及单元测试咨询

问题分析与解决方案

核心配置误区

你的配置存在两个关键问题,导致OnMissingInstance未按预期触发:

  1. InsertOnInitial = true与OnMissingInstance互斥
    当设置InsertOnInitial = true时,MassTransit会在找不到对应Saga实例时自动创建新实例并进入Initially状态,完全跳过OnMissingInstance逻辑。这是你发送无效CorrelationId时未触发错误响应的直接原因。

  2. OnMissingInstance不处理已完成的Saga实例
    已进入终态的Saga实例仍然存在于存储中,因此不会触发OnMissingInstance。这种场景需要通过OnExists钩子检查实例状态并处理。

修正后的配置代码

Event(() => PaymentAuthorization, x => 
{
    x.CorrelateById(context => context.Message.AuthorizationRequestId);
    
    // 处理找不到实例的场景
    x.OnMissingInstance(async context =>
    {
        await context.RespondAsync(new AuthorizationError(
            context.CorrelationId, 
            "Saga instance not found"
        ));
    });
    
    // 移除InsertOnInitial = true,避免自动创建实例
    x.SetSagaFactory(ctx => new PaymentAuthorizationState
    {
        CorrelationId = ctx.Message.AuthorizationRequestId,
    });
    
    // 处理已存在但已完成的实例
    x.OnExists(async context =>
    {
        if (context.Saga.IsCompleted)
        {
            await context.RespondAsync(new AuthorizationError(
                context.CorrelationId, 
                "Saga instance has already completed"
            ));
            return false; // 返回false表示不继续处理该事件
        }
        return true; // 返回true表示正常处理事件
    });
});

Initially(
    When(PaymentAuthorization)
        .Activity(x => x.OfType<AuthorizationRequestedActivity>())
        .ThenAsync(async context =>
        {
            context.Saga.Status = PaymentAuthorizationStatus.Waiting;
            await context.RespondAsync(new PaymentAuthorizationResult(xxxx));
        })
        .Schedule(PaymentSessionExpired, ctx => new PaymentSessionExpired(ctx.Saga.CorrelationId))
        .TransitionTo(Waiting)
        .Catch<Exception>(binder => binder
            .ThenAsync(async context =>
            {
                await context.RespondAsync(new AuthorizationError(
                    context.Message.AuthorizationRequestId,  
                    "Error processing authorization"
                ));
            })
        ));

SetCompletedWhenFinalized();

OnMissingInstance单元测试示例

使用MassTransit内存测试框架编写验证用例:

[TestFixture]
public class PaymentAuthorizationSagaTests
{
    [Test]
    public async Task Invalid_CorrelationId_Triggers_Missing_Instance_Response()
    {
        var harness = new InMemoryTestHarness();
        var sagaHarness = harness.Saga<PaymentAuthorizationState, PaymentAuthorizationStateMachine>();

        await harness.Start();
        try
        {
            var requestId = Guid.NewGuid();
            var response = await harness.RequestClient<PaymentAuthorization>()
                .GetResponse<AuthorizationError>(new PaymentAuthorization 
                { 
                    AuthorizationRequestId = requestId 
                });

            Assert.Multiple(() =>
            {
                Assert.That(response.Message.CorrelationId, Is.EqualTo(requestId));
                Assert.That(response.Message.Message, Is.EqualTo("Saga instance not found"));
                Assert.That(sagaHarness.Created.Contains(requestId), Is.False);
            });
        }
        finally
        {
            await harness.Stop();
        }
    }

    [Test]
    public async Task Completed_Saga_Triggers_Exists_Error_Response()
    {
        var harness = new InMemoryTestHarness();
        var sagaHarness = harness.Saga<PaymentAuthorizationState, PaymentAuthorizationStateMachine>();

        await harness.Start();
        try
        {
            var requestId = Guid.NewGuid();
            // 先触发Saga完成流程
            await harness.Bus.Publish(new PaymentAuthorization { AuthorizationRequestId = requestId });
            await harness.Bus.Publish(new PaymentCompleted { CorrelationId = requestId });

            // 等待Saga进入完成状态
            await sagaHarness.Consumed.Any<PaymentCompleted>(x => x.Context.CorrelationId == requestId);
            Assert.That(sagaHarness.Sagas[requestId].IsCompleted, Is.True);

            // 再次发送PaymentAuthorization事件
            var response = await harness.RequestClient<PaymentAuthorization>()
                .GetResponse<AuthorizationError>(new PaymentAuthorization 
                { 
                    AuthorizationRequestId = requestId 
                });

            Assert.Multiple(() =>
            {
                Assert.That(response.Message.CorrelationId, Is.EqualTo(requestId));
                Assert.That(response.Message.Message, Is.EqualTo("Saga instance has already completed"));
            });
        }
        finally
        {
            await harness.Stop();
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:05:54