如何通过队列实现线程复用与线程双向通信处理批量任务
问题根源
- 工作线程没有外层循环,执行完单次任务后线程就直接终止,自然不会复用拉取下一批任务
ManualResetEvent调用Set()后一直处于已触发状态,没有做重置,无法实现批次级的触发控制- 缺少双向同步信号:没有统计当前批次任务完成数的计数器,也没有通知控制线程批次已完成的信号量
- 你把整个工作线程的逻辑都包在
lock里,相当于所有线程串行执行任务,完全没有发挥多线程的作用,性能极差
修复后的实现代码
首先定义需要的同步变量:
Queue<Action> queueComputersToScan = new Queue<Action>(); // 两个同步事件:通知工作线程有任务、通知控制线程批次完成 ManualResetEvent workerEvent = new ManualResetEvent(false); ManualResetEvent controllerEvent = new ManualResetEvent(false); object syncLock = new object(); int batchsize = 20; int pendingTasksInBatch = 0; // 当前批次剩余未完成的任务数 List<Thread> threadsPerComputer = new List<Thread>();
控制线程代码:
Thread controllerThread = new Thread(() => { var allComputers = GetListOfComputers(); int totalBatches = (int)Math.Ceiling((decimal)allComputers.Count / batchsize); for (int x = 0; x < totalBatches; x++) { // 拉取当前批次的任务入队 var computers = allComputers.Skip(x * batchsize).Take(batchsize).ToList(); lock (syncLock) { foreach (var computer in computers) { // 修复闭包捕获问题,单独声明当前迭代变量 var currentComputer = computer; queueComputersToScan.Enqueue(() => ScanComputer(currentComputer)); } pendingTasksInBatch = computers.Count; } // 通知工作线程批次任务已准备好 workerEvent.Set(); // 等待当前批次所有任务执行完成 controllerEvent.WaitOne(); // 重置信号,为下一批次做准备 controllerEvent.Reset(); workerEvent.Reset(); } // 所有批次执行完成,通知工作线程退出 workerEvent.Set(); }); controllerThread.Start();
工作线程代码:
int maxThread = 20; for (int i = 0; i < maxThread; i++) { Thread thread = new Thread(() => { while (true) { // 等待工作信号 workerEvent.WaitOne(); Action currentJob = null; lock (syncLock) { // 队列还有任务就取任务 if (queueComputersToScan.Count > 0) { currentJob = queueComputersToScan.Dequeue(); } else { // 队列空了说明所有任务都结束了,退出循环 break; } } // 任务执行不放在lock块内,保证多线程并行 currentJob?.Invoke(); // 任务完成后更新剩余计数 lock (syncLock) { pendingTasksInBatch--; if (pendingTasksInBatch == 0) { // 批次所有任务完成,通知控制线程 controllerEvent.Set(); } } } }); thread.Name = $"ScanThread_{i}"; thread.Start(); Console.WriteLine($"{thread.Name.ToLower()} has started"); threadsPerComputer.Add(thread); } // 所有任务执行完成后等待工作线程退出 foreach (var thread in threadsPerComputer) { thread.Join(); }
核心修改点说明
- 给工作线程加了外层
while循环,实现线程复用,会一直拉取任务直到所有批次完成 - 新增了
controllerEvent作为反向通知信号,实现控制线程和工作线程的双向同步 - 增加了
pendingTasksInBatch计数器,用来统计当前批次剩余未完成任务,计数为0时通知控制线程拉取下一批 - 把任务执行逻辑移出了
lock块,只有队列操作和计数修改加锁,保证多线程并行执行任务 - 修复了lambda闭包捕获循环变量的问题,避免所有任务都用最后一个computer变量的BUG
- 每个批次完成后都会重置两个事件的状态,保证下一批次的同步逻辑正常运行
内容的提问来源于stack exchange,提问作者MiralShah
相关产品推荐
相关产品推荐

