MassTransit处理Job抛出TaskCanceledException异常问题求助
问题根因
异常触发的核心是JobService协调作业时找不到对应作业的Saga状态实例,删除Saga操作触发任务取消异常,问题来自两处配置错误:
IJobConsumer实现类不允许手动绑定到自定义接收端点,必须由ServiceInstance自动注册- 作业消息不能直接通过
SendEndpoint发送到自定义队列,必须使用MassTransit提供的IJobPublisher发布
修复步骤
1. 修正消费者端配置
删除手动声明的loan-request-processing接收端点配置,调整后的消费者配置如下:
services.AddMassTransit(x => { x.AddDelayedMessageScheduler(); x.AddConsumer<LoanRequestJobConsumer>(cfg => { cfg.Options<JobOptions<LoanRequestBroker>>(options => { options.SetJobTimeout(TimeSpan.FromMinutes(5)); options.SetConcurrentJobLimit(10); }); }); x.SetKebabCaseEndpointNameFormatter(); x.UsingRabbitMq((context, cfg) => { cfg.UseDelayedMessageScheduler(); cfg.ServiceInstance(instance => { instance.ConfigureJobServiceEndpoints(js => { js.SagaPartitionCount = 1; js.FinalizeCompleted = true; }); // 自动注册所有JobConsumer对应的端点,不要手动配置接收端点 instance.ConfigureEndpoints(context); }); }); }); services.AddMassTransitHostedService();
2. 修正生产者端配置和消息发布逻辑
首先调整生产者的MassTransit配置,添加作业客户端支持:
services.AddMassTransit(x => { x.SetKebabCaseEndpointNameFormatter(); // 注册作业状态机支持 x.AddJobSagaStateMachines(); x.UsingRabbitMq((context, cfg) => { cfg.ConfigureEndpoints(context); }); }); services.AddMassTransitHostedService();
然后修改消息入队代码,通过IJobPublisher发布作业:
// 构造函数注入IJobPublisher private readonly IJobPublisher<LoanRequestBroker> _jobPublisher; // 发布作业消息 var jobId = await _jobPublisher.Publish(loanRequest.Adapt<LoanRequestBroker>());
3. 生产环境优化(可选)
默认的内存Saga存储仅适合本地开发使用,生产环境建议替换为Redis、PostgreSQL等持久化存储,避免服务重启丢失作业状态。
内容的提问来源于stack exchange,提问作者Nikita Ivanov
相关产品推荐
相关产品推荐

