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

如何通过队列实现线程复用与线程双向通信处理批量任务

问题根源
  • 工作线程没有外层循环,执行完单次任务后线程就直接终止,自然不会复用拉取下一批任务
  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 17:15:00