MassTransit Saga功能报错“Send saga is no longer available in IoC”的解决方案求助
MassTransit Saga功能报错“Send saga is no longer available in IoC”的解决方案求助
大家好,我在使用MassTransit的Saga功能时遇到了一个棘手的问题,完全搞不懂下面这个错误的原因,想请各位大佬帮忙分析一下~
我遇到的错误信息如下:
Error: MassTransit.ReceiveTransport[0] R-FAULT rabbitmq://localhost/monitoring-job-saga 489a0000-4100-0250-1852-08dd500596e7 Shared.MessagingContracts.JobSubmitted MainApp.MonitoringJobState(00:00:00.0334249) System.NotSupportedException: Send saga is no longer available in IoC at MassTransit.Configuration.RegistrationServiceCollectionExtensions.TempSagaRepository`1.Send[T](ConsumeContext`1 context, ISagaPolicy`2 policy, IPipe`1 next) in /_/src/MassTransit/DependencyInjection/Configuration/RegistrationServiceCollectionExtensions.cs:Line 93 at MassTransit.Middleware.CorrelatedSagaFilter`2.Send(ConsumeContext`1 context, IPipe`1 next)
下面是我的MassTransit配置代码:
services.AddMassTransit(x => { x.AddSagaStateMachine<MonitoringJobStateMachine, MonitoringJobState>() .InMemoryRepository(); x.UsingRabbitMq((context, cfg) => { cfg.Host(new Uri(rabbitHost), h => { h.Username(rabbitUser); h.Password(rabbitPass); }); cfg.UseDelayedMessageScheduler(); cfg.ReceiveEndpoint(sagaQueue, e => { e.StateMachineSaga( context.GetRequiredService<MonitoringJobStateMachine>(), context.GetRequiredService<ISagaRepository<MonitoringJobState>>()); }); }); });
还有我的MonitoringJobStateMachine类代码:
using MassTransit; using MainApp.Data; using Shared.MessagingContracts; namespace MainApp { public class MonitoringJobStateMachine : MassTransitStateMachine<MonitoringJobState> { public State Submitted { get; private set; } = null!; public State Processing { get; private set; } = null!; public State Completed { get; private set; } = null!; public State Failed { get; private set; } = null!; public Event<JobSubmitted> JobSubmittedEvent { get; private set; } = default!; public Event<JobCompleted> JobCompletedEvent { get; private set; } = default!; public Event<JobFailed> JobFailedEvent { get; private set; } = default!; public Schedule<MonitoringJobState, JobTimeout> JobTimeoutSchedule { get; private set; } = default!; public MonitoringJobStateMachine() { InstanceState(x => x.CurrentState); Event(() => JobSubmittedEvent, x => { x.CorrelateById(ctx => ctx.Message.CorrelationId); x.InsertOnInitial = true; }); Event(() => JobCompletedEvent, x => { x.CorrelateById(ctx => ctx.Message.CorrelationId); }); Event(() => JobFailedEvent, x => { x.CorrelateById(ctx => ctx.Message.CorrelationId); }); // 配置超时调度器 Schedule(() => JobTimeoutSchedule, x => x.TimeoutTokenId, s => { s.Delay = TimeSpan.FromSeconds(10); s.Received = r => r.CorrelateById(ctx => ctx.Message.CorrelationId); }); // 初始流转:JobSubmitted -> 发送ProcessJob并计划超时 Initially( When(JobSubmittedEvent) .ThenAsync(async context => { // 从入站消息中获取值 context.Saga.SubmittedAt = context.Message.Timestamp; context.Saga.Regions = context.Message.Regions; context.Saga.CurrentAttempt = 0; if (context.Saga.Regions != null && context.Saga.Regions.Count > 0) context.Saga.CurrentRegion = context.Saga.Regions[0]; Console.WriteLine($"[Saga] Job {context.Saga.CorrelationId} submitted. Starting in Region {context.Saga.CurrentRegion}."); // 向区域专属队列发送ProcessJob命令(例如"jobs-de") await context.Publish<ProcessJob>(new ProcessJobCommand { CorrelationId = context.Saga.CorrelationId, Region = context.Saga.CurrentRegion, Attempt = context.Saga.CurrentAttempt }, publishContext => { // 显式设置消息目标地址: publishContext.DestinationAddress = new Uri($"queue:jobs-{context.Saga.CurrentRegion.ToLower()}"); }); }) // 在流转链中直接计划超时 .Schedule(JobTimeoutSchedule, ctx => new JobTimeoutMessage { CorrelationId = ctx.Saga.CorrelationId }, ctx => TimeSpan.FromSeconds(10)) .TransitionTo(Processing) ); During(Processing, When(JobCompletedEvent) .Then(ctx => { Console.WriteLine($"[Saga] Job {ctx.Saga.CorrelationId} completed successfully in region {ctx.Message.Region}."); }) .Unschedule(JobTimeoutSchedule) .TransitionTo(Completed), When(JobFailedEvent) .ThenAsync(async ctx => { Console.WriteLine($"[Saga] Job {ctx.Saga.CorrelationId} failed in region {ctx.Message.Region}. Error: {ctx.Message.ErrorMessage}"); // 增加重试次数 ctx.Saga.CurrentAttempt++; if (ctx.Saga.Regions != null && ctx.Saga.CurrentAttempt < ctx.Saga.Regions.Count) { // 设置新区域 ctx.Saga.CurrentRegion = ctx.Saga.Regions[ctx.Saga.CurrentAttempt]; Console.WriteLine($"[Saga] Fallback: Retrying job {ctx.Saga.CorrelationId} in region {ctx.Saga.CurrentRegion}."); // 向新的区域队列发送新的ProcessJob命令 await ctx.Publish<ProcessJob>(new ProcessJobCommand { CorrelationId = ctx.Saga.CorrelationId, Region = ctx.Saga.CurrentRegion, Attempt = ctx.Saga.CurrentAttempt }, publishContext => { publishContext.DestinationAddress = new Uri($"queue:jobs-{ctx.Saga.CurrentRegion.ToLower()}"); }); } }) // 失败(或超时)后计划新的超时 .Schedule(JobTimeoutSchedule, ctx => new JobTimeoutMessage { CorrelationId = ctx.Saga.CorrelationId }, ctx => TimeSpan.FromSeconds(10)) .IfElse(ctx => ctx.Saga.Regions != null && ctx.Saga.CurrentAttempt < ctx.Saga.Regions.Count, binder => binder.TransitionTo(Processing), binder => binder.TransitionTo(Failed) ), When(JobTimeoutSchedule.Received) .ThenAsync(async ctx => { Console.WriteLine($"[Saga] Timeout in region {ctx.Saga.CurrentRegion} for job {ctx.Saga.CorrelationId}."); // 超时后同样增加重试次数并启动 fallback ctx.Saga.CurrentAttempt++; if (ctx.Saga.Regions != null && ctx.Saga.CurrentAttempt < ctx.Saga.Regions.Count) { ctx.Saga.CurrentRegion = ctx.Saga.Regions[ctx.Saga.CurrentAttempt]; Console.WriteLine($"[Saga] Fallback after timeout: Retrying job {ctx.Saga.CorrelationId} in region {ctx.Saga.CurrentRegion}."); await ctx.Publish<ProcessJob>(new ProcessJobCommand { CorrelationId = ctx.Saga.CorrelationId, Region = ctx.Saga.CurrentRegion, Attempt = ctx.Saga.CurrentAttempt }, publishContext => { publishContext.DestinationAddress = new Uri($"queue:jobs-{ctx.Saga.CurrentRegion.ToLower()}"); }); } }) .Schedule(JobTimeoutSchedule, ctx => new JobTimeoutMessage { CorrelationId = ctx.Saga.CorrelationId }, ctx => TimeSpan.FromSeconds(10)) ); } } }
备注:内容来源于stack exchange,提问作者LUCKYONE
相关产品推荐
相关产品推荐

