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

如何在C#中实现支持Worker重入队的有限并发工作队列?

C# 实现带动态入队的工作队列与Worker协调终止方案

需求描述

  • 可预先填充初始工作项的工作队列
  • 初始化后启动若干Worker任务/线程从队列取工作项
  • 每个Worker处理工作项时,允许将新工作项重新入队
  • 所有Worker持续运行,直到同时满足两个条件:
    1. 当前队列无剩余工作项
    2. 没有Worker正在处理工作项(不会再有新工作项入队)

初始实现与问题

初始代码

var queue = new ConcurrentQueue<WorkItem>();
// ... 向队列填充初始工作项 ...

await Task.WhenAll(Enumerable.Range(0, nWorkers).Select(_ => Task.Run(async () => {
  while(queue.TryDequeue(out var item))
  {
     // 异步处理工作,可能重新入队等操作
  }
}).ToArray());

存在问题

当某个Worker因TryDequeue失败(队列暂时为空)而提前终止时,其他活跃Worker仍可能处理任务并重新入队新工作项,导致后续入队的任务无人处理。

最简实现方案

核心思路是用原子计数器跟踪正在处理的任务数,结合队列状态判断是否可以终止Worker循环。以下是两种不同场景的实现:

场景1:轻量任务(简单实现)

适合任务量不大、对CPU消耗不敏感的场景,用延迟重试避免空转:

using System.Threading;

public class WorkItem { /* 自定义工作项结构 */ }

var queue = new ConcurrentQueue<WorkItem>();
// 填充初始工作项
// queue.Enqueue(new WorkItem());

int processingCount = 0;
int workerCount = 4; // 自定义Worker数量

var workers = Enumerable.Range(0, workerCount).Select(_ => Task.Run(async () => {
    while (true)
    {
        if (queue.TryDequeue(out var item))
        {
            // 开始处理,原子递增计数
            Interlocked.Increment(ref processingCount);
            try
            {
                // 异步处理工作项,可在此调用queue.Enqueue添加新任务
                // await ProcessWorkItem(item);
            }
            finally
            {
                // 处理完成,原子递减计数
                Interlocked.Decrement(ref processingCount);
            }
        }
        else
        {
            // 队列为空时,检查是否无活跃处理任务
            if (Volatile.Read(ref processingCount) == 0)
            {
                break; // 满足终止条件,退出循环
            }
            // 短暂等待后重试,避免CPU空转
            await Task.Delay(10);
        }
    }
})).ToArray();

await Task.WhenAll(workers);

场景2:高吞吐任务(高效实现)

用SemaphoreSlim替代延迟重试,减少不必要的等待,提升吞吐量:

using System.Threading;

public class WorkItem { /* 自定义工作项结构 */ }

var queue = new ConcurrentQueue<WorkItem>();
var semaphore = new SemaphoreSlim(0);
int processingCount = 0;
int workerCount = 4;

// 填充初始任务并释放对应信号量
var initialItems = new List<WorkItem>(); // 初始工作项集合
foreach (var item in initialItems)
{
    queue.Enqueue(item);
    semaphore.Release();
}

var workers = Enumerable.Range(0, workerCount).Select(_ => Task.Run(async () => {
    while (true)
    {
        await semaphore.WaitAsync();
        
        if (queue.TryDequeue(out var item))
        {
            Interlocked.Increment(ref processingCount);
            try
            {
                // 处理任务,若生成新任务则入队并释放信号
                // var newTasks = await ProcessWorkItem(item);
                // foreach (var task in newTasks)
                // {
                //     queue.Enqueue(task);
                //     semaphore.Release();
                // }
            }
            finally
            {
                Interlocked.Decrement(ref processingCount);
            }
        }
        
        // 检查终止条件:队列为空且无活跃处理任务
        if (queue.IsEmpty && Volatile.Read(ref processingCount) == 0)
        {
            // 释放所有等待的Worker,让它们同步退出
            semaphore.Release(workerCount);
            break;
        }
    }
})).ToArray();

await Task.WhenAll(workers);
semaphore.Dispose();

关键细节说明

  • 原子计数:用Interlocked确保处理任务数的增减是线程安全的
  • Volatile读取:保证读取的processingCount是最新值,避免线程缓存导致的判断错误
  • 信号量唤醒:高吞吐场景下,用信号量替代延迟,Worker仅在有新任务时被唤醒,减少资源浪费

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 00:11:02