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

如何均衡分配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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 05:50:19