如何基于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
相关产品推荐
相关产品推荐

