MassTransit Saga OnMissingInstance配置异常及单元测试咨询
问题分析与解决方案
核心配置误区
你的配置存在两个关键问题,导致OnMissingInstance未按预期触发:
InsertOnInitial = true与OnMissingInstance互斥
当设置InsertOnInitial = true时,MassTransit会在找不到对应Saga实例时自动创建新实例并进入Initially状态,完全跳过OnMissingInstance逻辑。这是你发送无效CorrelationId时未触发错误响应的直接原因。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
相关产品推荐
相关产品推荐

