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

MassTransit配置疑问:如何设置作业消费者单次仅处理一个作业?

问题分析与解决

问题描述

发布消息启动作业时,已配置消费者为单次仅处理一个作业,其余作业应在队列等待当前作业完成。但当一个作业启动后,再发布更多作业时,新作业会被发送到job_error队列,错误信息为:

The JobSubmitted event is not handled during the Started state for the JobStateMachine state machine

程序配置

builder.Services.AddMassTransit(x =>
{
    x.AddDelayedMessageScheduler();
    x.SetEndpointNameFormatter(new SnakeCaseEndpointNameFormatter(includeNamespace: true));

    x.AddConsumer<OnArchiveTimedRequestsToCheaperStorage, OnArchiveTimedRequestsToCheaperStorageDefinition>();

    //x.AddConsumers(typeof(Program).Assembly);

    x.SetJobConsumerOptions();
    x.SetInMemorySagaRepositoryProvider();
    x.AddJobSagaStateMachines();

    x.UsingAzureServiceBus((context, cfg) =>
    {
        cfg.Host(builder.Configuration["AzureServiceBusConfiguration:ConnectionString"]);
        cfg.SetNamespaceSeparatorTo("_");

        cfg.UseServiceBusMessageScheduler();
        cfg.ConfigureEndpoints(context);
    });
});

消费者代码

public class OnArchiveTimedRequestsToCheaperStorage : IJobConsumer<ArchiveTimedRequestsToCheaperStorage>
{
    public async Task Run(JobContext<ArchiveTimedRequestsToCheaperStorage> context)
    {
        bool run = true;
        while (run)
        {
            await Task.Delay(2000);
        }
    }
}

public class OnArchiveTimedRequestsToCheaperStorageDefinition : ConsumerDefinition<OnArchiveTimedRequestsToCheaperStorage>
{
    protected override void ConfigureConsumer(IReceiveEndpointConfigurator endpointConfigurator
        , IConsumerConfigurator<OnArchiveTimedRequestsToCheaperStorage> consumerConfigurator
        , IRegistrationContext context)
    {
        consumerConfigurator.Options<JobOptions<ArchiveTimedRequestsToCheaperStorage>>(options => options
            .SetConcurrentJobLimit(1));
    }
}

Azure Service Bus 队列截图

Azure Service Bus 队列截图

解决步骤

1. 提交作业使用唯一JobId

错误核心原因是重复使用同一个JobId提交作业,导致同一个JobStateMachine Saga实例收到JobSubmitted事件,但此时该实例处于Started状态,无法处理该事件。提交作业时必须生成唯一JobId:

var endpoint = await bus.GetSendEndpoint(new Uri($"queue:{KebabCaseEndpointNameFormatter.Instance.SanitizeName(typeof(ArchiveTimedRequestsToCheaperStorage).Name)}"));

await endpoint.Send<ArchiveTimedRequestsToCheaperStorage>(new {}, context =>
{
    // 每次提交生成新的JobId
    context.SetJobId(Guid.NewGuid());
});

2. 移除重复配置项

删除UsingAzureServiceBus配置块内的x.AddDelayedMessageScheduler();,该配置已在全局添加,重复添加会导致调度器冲突。

3. 调整Azure Service Bus配置顺序

确保UseServiceBusMessageScheduler()在ConfigureEndpoints(context)之前调用,让所有端点正确使用Service Bus消息调度器,实现作业等待队列逻辑。

4. 验证并发限制配置

消费者定义中SetConcurrentJobLimit(1)的配置是正确的,该设置会让超出并发数的作业进入延迟队列,等待当前作业完成后再执行。

内容的提问来源于stack exchange,提问作者Brandon

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 12:02:11