【问题标题】:C#: limit maximum of concurrent operation with Parallel.ForEach and async ActionC#:使用 Parallel.ForEach 和异步操作限制最大并发操作
【发布时间】:2018-09-14 06:16:38
【问题描述】:

我正在尝试使用 asp.net core 2.1 实现自托管 Web 服务,但遇到了实现后台长时间执行任务的问题。

由于每个ProcessSingle方法(在下面的代码sn-p中)的高CPU负载和时间消耗,我想限制同时执行的任务的数量。但是我几乎可以立即看到Parallel.ForEachstart 中的所有任务,尽管我设置了MaxDegreeOfParallelism = 3

我的代码是(简化版):

public static async Task<int> Work()
{
    var id = await CreateIdInDB() // async create record in DB

    // run background task, don't wait when it finishes
    Task.Factory.StartNew(async () => {
        Parallel.ForEach(
            listOfData,
            new ParallelOptions { CancellationToken = token, MaxDegreeOfParallelism = 3 },
            async x => await ProcessSingle(x));
    });

    // return created id immediately
    return id;
}

public static async Task ProcessSingle(MyInputData inputData)
{
    var dbData = await GetDataFromDb(); // get data from DB async using Dapper
    // some lasting processing (sync)
    await SaveDataToDb(); // async save processed data to DB using Dapper
}

如果我理解正确,问题出在 Parallel.ForEach 内的async x =&gt; await ProcessSingle(x),不是吗?

有人可以描述一下,它应该如何以正确的方式实施?

更新

由于我的问题有些模棱两可,有必要关注主要方面:

  1. ProcessSingle方法分三部分:

    • 从数据库异步获取数据

    • 进行长时间高 CPU 负载的数学计算

    • 将结果保存到数据库异步

  2. 这个问题包括两个独立的:

    • 如何降低 CPU 使用率(例如同时运行不超过三个数学计算)?

    • 如何保持 ProcessSingle 方法的结构 - 因为异步 DB 调用而使它们保持异步。

希望现在会更清楚。

附:已经给出了合适的答案,它可以工作(特别感谢@MatrixTai)。编写此更新是为了进行一般说明。

【问题讨论】:

标签: c# async-await task task-parallel-library


【解决方案1】:

更新

正如我刚刚注意到你在评论中提到的,问题是由数学计算引起的。

计算和更新DB的部分最好分开。

计算部分使用Parallel.ForEach(),这样可以优化你的工作,你可以控制线程数。

只有在所有这些任务完成之后。使用async-await 将您的数据更新到数据库,而不需要我提到的SemaphoreSlim。

public static async Task<int> Work()
{
    var id = await CreateIdInDB() // async create record in DB

    // run background task, don't wait when it finishes
    Task.Run(async () => {

        //Calculation Part
        ConcurrentBag<int> data = new ConcurrentBag<int>();
        Parallel.ForEach(
            listOfData,
            new ParallelOptions { CancellationToken = token, MaxDegreeOfParallelism = 3 },
            x => {ConcurrentBag.Add(calculationPart(x))});

        //Update DB part
        int[] data_arr = data.ToArray();
        List<Task> worker = new List<Task>();
        foreach (var i in data_arr)
        {
            worker.Add(DBPart(x));
        }
        await Task.WhenAll(worker);
    });

    // return created id immediately
    return id;
}

当您在Parallel.forEach 中使用async-await 时,它们肯定是一起开始的。

首先,请阅读question 第一个和第二个答案。将这两者结合起来毫无意义。

其实async-await会最大化可用线程的使用率,所以直接用吧。

public static async Task<int> Work()
{
    var id = await CreateIdInDB() // async create record in DB

    // run background task, don't wait when it finishes
    Task.Run(async () => {
        List<Task> worker = new List<Task>();
        foreach (var i in listOfData)
        {
            worker.Add(ProcessSingle(x));
        }
        await Task.WhenAll(worker);
    });

    // return created id immediately
    return id;
}

但是这里有另一个问题,在这种情况下,这些任务仍然一起开始,消耗你的 CPU 使用率。

为了避免这种情况,请使用SemaphoreSlim

public static async Task<int> Work()
{
    var id = await CreateIdInDB() // async create record in DB

    // run background task, don't wait when it finishes
    Task.Run(async () => {
        List<Task> worker = new List<Task>();
        //To limit the number of Task started.
        var throttler = new SemaphoreSlim(initialCount: 20);
        foreach (var i in listOfData)
        {
            await throttler.WaitAsync();
            worker.Add(Task.Run(async () =>
            {
                await ProcessSingle(x);
                throttler.Release();
            }
            ));
        }
        await Task.WhenAll(worker);
    });

    // return created id immediately
    return id;
}

阅读更多How to limit the amount of concurrent async I/O operations?。

另外,当简单的Task.Run() 足以完成你想做的工作时,不要使用Task.Factory.StartNew(),请阅读Stephen Cleary 的这篇出色的article。

【讨论】:

  • 有一个内置类用于与 DOP 并行执行作业,它是 ActionBlock。在任何情况下,并行执行多个慢速数据库查询更有可能降低性能
  • @PanagiotisKanavos,我同意如果问题与数据库查询有关,则更有可能降低性能。为自己完成这项任务,我更有可能同步完成。
  • @MatrixTai,不,不,数据库通信完全没有问题:) 你已经回答了我想要的。谢谢)
  • @user1820686 ,等等,我刚刚注意到您在评论中写道,问题是由数学计算引起的。然后,它主要是 CPU-Tasked。在这种情况下,最好使用 PMF 方法,通过Parallel.forEach 使其同步和dun,最后收集async-await 写入数据库的所有数据。我会更新我的答案。
  • @user1820686 你的实际问题是什么? DOP 为 3 的 Parallel.ForEach 将同时运行 3 个任务。它将数据拆分为 3 种方式,并将每个分区传递给一个 single 任务。您的问题的 错误 是通过使 ProcessSingle 异步,您为每个单独的数据项启动 另一个 任务
【解决方案2】:

如果您更熟悉“传统”并行处理概念,请像这样重写您的 ProcessSingle() 方法:

public static void ProcessSingle(MyInputData inputData)
{
    var dbData = GetDataFromDb(); // get data from DB async using Dapper
    // some lasting processing (sync)
    SaveDataToDb(); // async save processed data to DB using Dapper
}

当然,您最好也以类似的方式更改 Work() 方法。

【讨论】:

  • 看起来很可疑,您会将Task 与Parallel.ForEach 迭代混合
  • @MickyD:是的,我会的。不过,这不是问题。混合Parallel.ForEach 和async / await 是。 Parallel.ForEach 本身使用任务。
  • 这是个问题。请参阅上面的链接。
  • @MickyD:对不起,哪个链接? stackoverflow.com/questions/11564506/… 没有提到任何关于混合 Tasks 和 Parallel.ForEach 的内容。这完全是关于混合任务和异步/等待。我根本不会使用任何 async/await。
  • 不正确。您已经表明您将使用Task。实际的await 是无关。 Read this。此外,即使您纯粹使用了Parallel.ForEach,将线程池线程与诸如GetDataFromDb() 和SaveDataToDb() 之类的阻塞I/O 任务一起使用也是对完美线程的浪费。您need 使用 IOCP,这在您的示例中没有。 async/await 非常适合 I/O,但不能与 Parallel.ForEach 一起使用。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-12-28
  • 2023-03-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多