如何均衡分配ConcurrentQueue任务至Worker并支持优先级入队?
任务分配与代码实现
优化后代码结构
public class WorkItem { public string Name { get; set; } // 标记是否为紧急任务 public bool IsUrgent { get; set; } } public class Worker { private readonly SemaphoreSlim _semaphore; // 允许的最大并发数 public int ConcurrentLimit { get; } // 当前正在运行的任务数 public int CurrentRunningCount => ConcurrentLimit - _semaphore.CurrentCount; public Worker(int concurrentLimit) { ConcurrentLimit = concurrentLimit; _semaphore = new SemaphoreSlim(concurrentLimit, concurrentLimit); } public async Task DoWork(WorkItem workItem) { await _semaphore.WaitAsync(); try { await Task.Delay(1000); // 模拟Microsoft Graph查询操作 } finally { _semaphore.Release(); } } } public class Engine { // 双队列实现紧急任务优先:先处理紧急队列,再处理普通队列 private readonly ConcurrentQueue<WorkItem> _urgentWorkItems = new(); private readonly ConcurrentQueue<WorkItem> _regularWorkItems = new(); private readonly List<Worker> _workers = new(); private readonly CancellationTokenSource _cts; public int ConcurrentThreadsForEachWorker { get; private set; } public Engine(int workers, int threads) { ConcurrentThreadsForEachWorker = threads; for (int i = 0; i < workers; i++) { _workers.Add(new Worker(threads)); } _cts = new CancellationTokenSource(); } // 添加普通任务到队列 public void EnqueueWorkItem(WorkItem workItem) { _regularWorkItems.Enqueue(workItem); } // 添加紧急任务到队列头部(优先处理) public void EnqueueUrgentWorkItem(WorkItem workItem) { workItem.IsUrgent = true; _urgentWorkItems.Enqueue(workItem); } public async Task RunAsync(CancellationToken token) { var linkedToken = CancellationTokenSource.CreateLinkedTokenSource(token, _cts.Token).Token; // 启动所有Worker的任务处理循环 var workerTasks = _workers.Select(worker => ProcessWorkerTasks(worker, linkedToken)).ToList(); await Task.WhenAll(workerTasks); } private async Task ProcessWorkerTasks(Worker worker, CancellationToken token) { while (!token.IsCancellationRequested) { WorkItem workItem = null; // 优先获取紧急任务 if (!_urgentWorkItems.TryDequeue(out workItem)) { // 无紧急任务时获取普通任务 _regularWorkItems.TryDequeue(out workItem); } if (workItem != null) { // 异步执行任务,不阻塞当前循环,保证Worker能继续拿新任务 _ = ProcessWorkItemAsync(worker, workItem, token); } else { // 无任务时短暂等待,避免CPU空转 await Task.Delay(100, token); } } } private async Task ProcessWorkItemAsync(Worker worker, WorkItem workItem, CancellationToken token) { try { await worker.DoWork(workItem); // 可选:打印任务完成日志 // Console.WriteLine($"WorkItem {workItem.Name} -> Worker {_workers.IndexOf(worker)+1} (已完成, 当前运行数: {worker.CurrentRunningCount})"); } catch (OperationCanceledException) { // 任务被取消,无需额外处理 } catch (Exception ex) { // 处理任务执行异常,比如记录日志或重试 // Console.WriteLine($"WorkItem {workItem.Name} 执行失败: {ex.Message}"); } } // 停止引擎,触发所有Worker停止处理任务 public void Stop() { _cts.Cancel(); } }
任务分配规则
当Engine构造函数传入workers=2、threads=4时,任务需按以下规则均衡分配:
WorkItem 1 -> Worker 1 (running 1)
WorkItem 2 -> Worker 2 (running 1)
WorkItem 3 -> Worker 1 (running 2)
WorkItem 4 -> Worker 2 (running 2)
WorkItem 5 -> Worker 1 (running 3)
WorkItem 6 -> Worker 2 (running 3)
WorkItem 7 -> Worker 1 (running 4 - full)
WorkItem 8 -> Worker 2 (running 4 - full)
WorkItem 9 -> 所有Worker已满,等待空闲
业务背景
需处理约1000个收件箱,Worker对应Microsoft Graph Client,WorkItem对应收件箱,每个Client可同时执行4个查询,需最大化处理速度。此外,Microsoft Graph的新邮件通知需实时入队处理,且需支持将紧急任务添加至队列头部优先处理。
实现要点
- Worker并发控制:给每个Worker添加
SemaphoreSlim,严格限制同时运行的任务数,避免单个Graph Client并发过高触发平台限流。 - 均衡任务分配:每个Worker独立循环从共享队列取任务,天然实现轮询式的均衡分配,符合要求的任务分配顺序。
- 紧急任务优先:通过双队列分离紧急和普通任务,Worker会优先处理紧急队列中的任务,满足新邮件实时处理和紧急任务插队的需求。
- 性能优化:无任务时Worker短暂等待后再重试,避免CPU空转浪费资源。
- 优雅停止:支持通过
CancellationToken触发引擎停止,所有正在执行的任务会被妥善收尾,不会强制中断。
内容的提问来源于stack exchange,提问作者Alireza Noori
相关产品推荐
相关产品推荐

