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

Task库并行任务优化求助:任务完成后自动从数据库获取新任务

并行任务调度与动态数据获取问题解决

我使用Task库并行运行多个任务,最初的实现代码如下:

Parallel.ForEach(DatafromDB,
  item => { DownloadSSR(mediator, sourcePath, item,stoppingToken).Wait();
});

但这段代码并未真正实现高效并行——比如3个任务中Task1耗时15分钟,其余两个各耗时1分钟时,仍需等待所有任务全部完成。

为解决并行问题,我尝试使用MaxDegreeOfParallelism参数调整:

Parallel.ForEach(DatafromDB,
  new ParallelOptions { MaxDegreeOfParallelism = 3}, 
  item => { DownloadSSR(mediator, sourcePath, item,stoppingToken).Wait();
});

现在任务可以并行运行,但存在新问题:仅能首次从数据库获取数据,比如取3条数据后,后续任务会重复使用该数据源,无法动态获取新数据。

我的核心需求是:

  • 并行运行任务
  • 单个任务完成后无需等待其他任务
  • 单个任务完成后,自动从数据库获取新数据并分配新任务

解决方案:异步任务池+动态数据拉取

Parallel.ForEach基于静态数据源迭代,无法满足动态拉取数据的需求。改用异步任务池+循环拉取数据的方式,实现任务完成后自动补新的调度逻辑:

示例代码

// 设置最大并行任务数
int maxParallelTasks = 3;
var workerTasks = new List<Task>();

// 启动指定数量的工作线程
for (int i = 0; i < maxParallelTasks; i++)
{
    workerTasks.Add(Task.Run(async () =>
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            // 从数据库拉取单条未处理的数据(需自行实现该方法,避免重复获取)
            var nextItem = await GetUnprocessedItemFromDBAsync(stoppingToken);
            
            // 无新数据则退出循环
            if (nextItem == null)
                break;
                
            // 异步执行下载任务,避免阻塞线程
            await DownloadSSR(mediator, sourcePath, nextItem, stoppingToken);
        }
    }, stoppingToken));
}

// 等待所有工作线程完成
await Task.WhenAll(workerTasks);

关键细节说明

  • GetUnprocessedItemFromDBAsync:需自行实现,通过数据库状态标记、分页或行锁机制,确保每次只拉取未处理的数据,避免重复分配
  • 全程异步:使用async/await替代.Wait(),避免线程阻塞,提升资源利用率
  • 取消机制:通过stoppingToken统一管理任务取消,收到停止信号时所有工作线程可优雅退出
  • 动态维持并行度:每个工作线程完成任务后立即拉取新数据,始终保持设定的最大并行数,不会因单个任务耗时久导致资源闲置

内容的提问来源于stack exchange,提问作者Abhijith Nayak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:53:21