.NET 8中MassTransit 8.25.0作业调度执行故障排查
问题:.NET 8 + MassTransit 8.25.0作业调度无法加入队列
背景
.NET 8环境下使用MassTransit 8.25.0开发包含调度器与作业消费者的Windows服务,遇到作业无法加入队列的问题。IPublishEndpoint为作用域服务,无法直接注入。
三次尝试及对应结果
尝试1:使用IClientFactory创建RequestClient提交作业
抛出错误:
MassTransit.EventExecutionException: The StartJobAttempt (Event) execution faulted MassTransit.PayloadNotFoundException: The payload was not found: MassTransit.MessageSchedulerContext
对应代码:
public class ImportScheduler( IScheduleConfig<ImportScheduler> config, ILogger<ImportScheduler> logger, IClientFactory clientFactory) : CronJobService(config.CronExpression, config.TimeZoneInfo, logger) { protected override async Task DoWork(CancellationToken cancellationToken) { var requestClient = clientFactory.CreateRequestClient<ImportRequested>(); var jobId = NewId.Next(); await requestClient.GetResponse<JobSubmissionAccepted>(new ImportRequested( jobId.ToGuid(), Guid.Parse("7718D173-1F3F-4F8C-B282-5B5A1C183BCE")), cancellationToken); } }
尝试2:通过IBus.GetPublishSendEndpoint创建ISendEndpoint提交作业
代码实现:
public class ImportScheduler( IScheduleConfig<ImportScheduler> config, ILogger<ImportScheduler> logger, IBus bus) : CronJobService(config.CronExpression, config.TimeZoneInfo, logger) { protected override async Task DoWork(CancellationToken cancellationToken) { var publishEndpoint = await bus.GetPublishSendEndpoint<SubmitJob<ImportRequested>>(); var jobId = NewId.Next(); await publishEndpoint.Send(new { JobId = jobId, Job = new ImportRequested( jobId.ToGuid(), Guid.Parse("7718D173-1F3F-4F8C-B282-5B5A1C183BCE")) }, cancellationToken); } }
尝试3:直接使用bus.Publish<SubmitJob>提交作业
出现与尝试1类似的PayloadNotFoundException错误。
额外遇到的空引用异常
System.NullReferenceException: Object reference not set to an instance of an object. at MassTransit.ExceptionInfoException..ctor(ExceptionInfo exceptionInfo) in /_/src/MassTransit.Abstractions/Exceptions/ExceptionInfoException.cs:line 14 at MassTransit.JobService.FinalizeJobConsumer`1.Consume(ConsumeContext`1 context) in /_/src/MassTransit/JobService/JobService/FinalizeJobConsumer.cs:line 43 at MassTransit.Middleware.MethodConsumerMessageFilter`2.MassTransit.IFilter<MassTransit.ConsumerConsumeContext<TConsumer,TMessage>>.Send(ConsumerConsumeContext`2 context, IPipe`1 next) in /_/src/MassTransit/Middleware/MethodConsumerMessageFilter.cs:line 28 at MassTransit.Configuration.PipeConfigurator`1.LastPipe.Send(TContext context) in /_/src/MassTransit.Abstractions/Middleware/Configuration/PipeBuilder.cs:line 123 at MassTransit.Consumer.DelegateConsumerFactory`1.Send[TMessage](ConsumeContext`1 context, IPipe`1 next) in /_/src/MassTransit/Consumers/Consumer/DelegateConsumerFactory.cs:line 29 at MassTransit.Consumer.DelegateConsumerFactory`1.Send[TMessage](ConsumeContext`1 context, IPipe`1 next) in /_/src/MassTransit/Consumers/Consumer/DelegateConsumerFactory.cs:line 39 at MassTransit.Middleware.ConsumerMessageFilter`2.MassTransit.IFilter<MassTransit.ConsumeContext<TMessage>>.Send(ConsumeContext`1 context, IPipe`1 next) in /_/src/MassTransit/Middleware/ConsumerMessageFilter.cs:line 48
当前MassTransit配置
services.AddMassTransit(x => { // some consumers x.AddConsumer<ImportConsumer>(cc => { cc.Options<JobOptions<ImportRequested>>(xx => xx.SetRetry(r => r.None()).SetJobTimeout(TimeSpan.FromHours(3)).SetConcurrentJobLimit(1)); }); x.SetJobConsumerOptions(); x.AddJobSagaStateMachines(); x.SetInMemorySagaRepositoryProvider(); x.SetKebabCaseEndpointNameFormatter(); x.UsingAzureServiceBus((context, cfg) => { cfg.Host(configuration.GetValue<string>("AzureServiceBus")); cfg.ConfigureEndpoints(context); }); });
特殊情况说明
若将消费者从IJobConsumer改为IConsumer,即可通过IBus.Publish正常将作业加入队列。
内容的提问来源于stack exchange,提问作者Yehor Androsov
相关产品推荐
相关产品推荐

