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

如何基于已执行线程数暂停与恢复Parallel.ForEach循环?

暂停/恢复并行任务的实现方案

首先明确:Parallel.ForEach本身不支持原生的暂停与恢复功能,它基于TPL的同步并行模型设计,会一次性调度受并发限制的任务,没有内置的任务暂停阻塞机制。而且你的现有代码存在两个关键问题:

  1. Parallel.ForEach无法正确处理async委托,会导致异步任务被“火并忘记”,循环会在异步任务完成前提前结束
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 07:45:35