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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.15 03:34:34