如何在C#中实现支持Worker重入队的有限并发工作队列?
C# 实现带动态入队的工作队列与Worker协调终止方案
需求描述
- 可预先填充初始工作项的工作队列
- 初始化后启动若干Worker任务/线程从队列取工作项
- 每个Worker处理工作项时,允许将新工作项重新入队
- 所有Worker持续运行,直到同时满足两个条件:
- 当前队列无剩余工作项
- 没有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
相关产品推荐
相关产品推荐

