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

如何基于Quartz.net实现生产者/消费者模式及任务依赖调度?

Quartz.NET 实现主任务等待子任务完成的队列式调度方案

核心思路

完全可以通过「主任务触发子任务+监听子任务完成状态+队列隔离子任务数据」的方式实现需求,核心逻辑是:主任务负责拆分数据、触发子任务,通过Quartz的监听机制或共享计数器等待所有子任务处理完各自队列后再结束,同时利用Job分组或自定义标识实现子任务的队列隔离。

具体实现步骤

1. 配置队列隔离与子任务标识

  • 给所有子任务分配统一的JobGroup(比如"SubJobs"),通过JobKey区分不同子任务对应的队列
  • 主任务通过JobDataMap给每个子任务传递专属的队列数据,子任务内部根据传入的数据处理对应队列

2. 主任务实现等待逻辑

  • 主任务触发子任务后,通过自定义JobListener监听子任务的完成事件
  • 用原子计数器记录未完成的子任务数量,所有子任务完成后触发主任务结束的信号

3. 子任务处理各自队列

  • 子任务从JobDataMap中读取专属队列数据,执行处理逻辑
  • 完成后由监听器更新计数器,触发主任务的完成信号

代码示例

主任务实现

[DisallowConcurrentExecution]
public class MainJob : IJob
{
    private readonly IScheduler _scheduler;
    private readonly ILogger<MainJob> _logger;

    public MainJob(IScheduler scheduler, ILogger<MainJob> logger)
    {
        _scheduler = scheduler;
        _logger = logger;
    }

    public async Task Execute(IJobExecutionContext context)
    {
        _logger.LogInformation("主任务启动,准备触发子任务");

        // 从已有数据拆分出两个子任务的队列数据
        var queue1Data = GetQueue1Data();
        var queue2Data = GetQueue2Data();

        // 初始化计数器(子任务数量)和完成信号
        var remainingSubJobs = new AtomicInteger(2);
        var allDoneSignal = new ManualResetEventSlim(false);

        // 注册子任务完成监听器
        var completionListener = new SubJobCompletionListener(remainingSubJobs, allDoneSignal);
        _scheduler.ListenerManager.AddJobListener(completionListener, GroupMatcher<JobKey>.GroupEquals("SubJobs"));

        try
        {
            // 触发子任务1
            var subJob1Key = new JobKey("SubJob_Queue1", "SubJobs");
            var subJob1Detail = JobBuilder.Create<SubJobQueue1>()
                .WithIdentity(subJob1Key)
                .UsingJobData("QueueData", queue1Data)
                .Build();

            await _scheduler.ScheduleJob(subJob1Detail, TriggerBuilder.Create().StartNow().Build());

            // 触发子任务2
            var subJob2Key = new JobKey("SubJob_Queue2", "SubJobs");
            var subJob2Detail = JobBuilder.Create<SubJobQueue2>()
                .WithIdentity(subJob2Key)
                .UsingJobData("QueueData", queue2Data)
                .Build();

            await _scheduler.ScheduleJob(subJob2Detail, TriggerBuilder.Create().StartNow().Build());

            _logger.LogInformation("所有子任务已触发,等待处理完成...");
            allDoneSignal.Wait();
            _logger.LogInformation("所有子任务处理完成,主任务结束");
        }
        finally
        {
            // 移除监听器,避免内存泄漏
            _scheduler.ListenerManager.RemoveJobListener(completionListener.Name);
        }
    }

    // 模拟从已有数据获取队列1的数据
    private List<string> GetQueue1Data() => new() { "Queue1-Data1", "Queue1-Data2" };
    // 模拟从已有数据获取队列2的数据
    private List<string> GetQueue2Data() => new() { "Queue2-Data1", "Queue2-Data2", "Queue2-Data3" };
}

子任务示例(以队列1为例)

public class SubJobQueue1 : IJob
{
    private readonly ILogger<SubJobQueue1> _logger;

    public SubJobQueue1(ILogger<SubJobQueue1> logger)
    {
        _logger = logger;
    }

    public async Task Execute(IJobExecutionContext context)
    {
        _logger.LogInformation("子任务1开始处理队列");
        var queueData = context.JobDetail.JobDataMap.Get("QueueData") as List<string>;

        foreach (var data in queueData)
        {
            // 模拟数据处理逻辑
            await Task.Delay(800);
            _logger.LogInformation("子任务1处理完成数据:{Data}", data);
        }

        _logger.LogInformation("子任务1队列处理完毕");
    }
}

子任务完成监听器

public class SubJobCompletionListener : IJobListener
{
    private readonly AtomicInteger _remainingJobs;
    private readonly ManualResetEventSlim _completionSignal;

    public SubJobCompletionListener(AtomicInteger remainingJobs, ManualResetEventSlim completionSignal)
    {
        _remainingJobs = remainingJobs;
        _completionSignal = completionSignal;
        Name = "SubJobCompletionListener";
    }

    public string Name { get; }

    public Task JobToBeExecuted(IJobExecutionContext context, CancellationToken cancellationToken = default)
        => Task.CompletedTask;

    public Task JobExecutionVetoed(IJobExecutionContext context, CancellationToken cancellationToken = default)
        => Task.CompletedTask;

    public Task JobWasExecuted(IJobExecutionContext context, JobExecutionException jobException, CancellationToken cancellationToken = default)
    {
        // 子任务完成,计数器减1
        var remaining = _remainingJobs.DecrementAndGet();
        if (remaining == 0)
        {
            _completionSignal.Set();
        }
        return Task.CompletedTask;
    }
}

注意事项

  • 因为主任务加了[DisallowConcurrentExecution],同一时间只会有一个主任务实例运行,避免了并发触发子任务的问题
  • 如果子任务可能执行失败,需要在监听器中处理异常场景(比如计数器重置、通知主任务终止)
  • 若需要更灵活的队列管理,可结合外部消息队列(如Redis队列):主任务将数据推入对应队列,子任务监听队列处理,主任务通过监听队列消费完成状态实现等待,解耦性更强

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 08:30:47