如何基于已执行线程数暂停与恢复Parallel.ForEach循环?
暂停/恢复并行任务的实现方案
首先明确:Parallel.ForEach本身不支持原生的暂停与恢复功能,它基于TPL的同步并行模型设计,会一次性调度受并发限制的任务,没有内置的任务暂停阻塞机制。而且你的现有代码存在两个关键问题:
- Parallel.ForEach无法正确处理
async委托,会导致异步任务被“火并忘记”,循环会在异步任务完成前提前结束 StatusofWorkitem.Add不是线程安全操作,多线程并发添加会引发竞争条件
针对你的需求(最大并发2、支持暂停/恢复、每个迭代发起API调用),推荐使用异步任务+并发控制信号量+暂停事件的方案,具体实现如下:
核心控制变量定义
// 控制暂停/恢复的事件,初始状态为允许执行 private readonly ManualResetEventSlim _pauseEvent = new ManualResetEventSlim(true); // 控制最大并发数为2的信号量 private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(2, 2); // 存储任务结果的线程安全集合(也可使用ConcurrentBag替代lock) private readonly List<WorkItem> _statusOfWorkitem = new List<WorkItem>(); private readonly object _lockObj = new object();
异步任务处理逻辑
public async Task ProcessWorkItemsAsync(IEnumerable<KeyValuePair<YourKeyType, YourValueType>> listOfWI) { var taskList = new List<Task>(); foreach (var currentWI in listOfWI) { // 等待暂停事件:如果处于暂停状态,会阻塞在此处直到恢复 _pauseEvent.Wait(); // 获取并发许可:最多同时有2个任务获取到许可 await _semaphore.WaitAsync(); // 启动异步任务处理当前项 taskList.Add(Task.Run(async () => { try { // 你的原有业务逻辑 ServiceInput objServiceInput = new ServiceInput { Login = newLoginDetails, Connectioninfo = newConnectionInfo, WorkItem = currentWI.Value, Services = mlserviceList }; string inputstring = JsonConvert.SerializeObject(objServiceInput, Newtonsoft.Json.Formatting.Indented); WIuploader newUpload = new WIuploader(logininfo, currentWI.Value, newConnectionInfo); Progress<WorkItem> report = new Progress<WorkItem>(workItemUpdate => { // 可在此添加进度更新逻辑 }); var workItem = await Task.Run(() => newUpload.Execute(report)); // 线程安全地添加结果 lock (_lockObj) { _statusOfWorkitem.Add(workItem); } } finally { // 无论任务成功/失败,都释放并发许可,让后续任务可以执行 _semaphore.Release(); } })); } // 等待所有任务执行完成 await Task.WhenAll(taskList); }
暂停与恢复方法
// 暂停任务:后续任务会阻塞在_pauseEvent.Wait()处,已执行的任务不受影响 public void PauseProcessing() { _pauseEvent.Reset(); } // 恢复任务:解除阻塞,后续任务继续执行 public void ResumeProcessing() { _pauseEvent.Set(); }
关键说明
- 并发控制:
SemaphoreSlim确保同时最多有2个API调用在执行,满足你的并发限制需求 - 暂停/恢复:
ManualResetEventSlim实现可重复的暂停与恢复操作,调用PauseProcessing()后,新的任务会停止获取并发许可,已启动的任务会继续完成;调用ResumeProcessing()后,阻塞的任务会继续执行 - 线程安全:通过
lock(或ConcurrentBag<WorkItem>)保证结果集合的线程安全 - 异步适配:完全基于
async/await模型,避免Parallel.ForEach处理异步任务的缺陷
额外注意事项
- 若需要支持任务取消,可额外引入
CancellationTokenSource,在任务中检查取消信号 - 使用完毕后需手动释放
ManualResetEventSlim和SemaphoreSlim资源,避免内存泄漏 Progress<WorkItem>的回调默认会捕获当前同步上下文(如UI线程),若在非UI环境下使用可忽略此特性
内容的提问来源于stack exchange,提问作者Natchatra
相关产品推荐
相关产品推荐

