【问题标题】:Parallel.ForEach faster than Task.WaitAll for I/O bound tasks?对于 I/O 绑定任务,Parallel.ForEach 比 Task.WaitAll 快吗?
【发布时间】:2019-09-25 16:46:08
【问题描述】:

我的程序有两个版本,它们向 Web 服务器提交约 3000 个 HTTP GET 请求。

第一个版本基于我阅读的 here。该解决方案对我来说很有意义,因为发出 Web 请求是受 I/O 限制的工作,并且将 async/await 与 Task.WhenAll 或 Task.WaitAll 一起使用意味着您可以一次提交 100 个请求,然后等待所有请求在提交接下来的 100 个请求之前完成,这样您就不会阻塞 Web 服务器。我很惊讶地看到这个版本在大约 12 分钟内完成了所有工作 - 比我预期的要慢。

第二个版本在 Parallel.ForEach 循环内提交所有 3000 个 HTTP GET 请求。我使用 .Result 等待每个请求完成,然后循环迭代中的其余逻辑才能执行。我认为这将是一个效率低得多的解决方案,因为使用线程并行执行任务通常更适合执行 CPU 密集型工作,但令我惊讶的是,这个版本在大约 3 分钟内完成了所有工作!

我的问题是为什么 Parallel.ForEach 版本更快?这是一个额外的惊喜,因为当我将相同的两种技术应用于不同 API/Web 服务器时,我的代码的版本 1 实际上比版本 2 快了大约 6分钟 - 这是我所期望的。两个不同版本的性能是否与 Web 服务器处理流量的方式有关?

您可以在下面看到我的代码的简化版本:

private async Task<ObjectDetails> TryDeserializeResponse(HttpResponseMessage response)
{
    try
    {
        using (Stream stream = await response.Content.ReadAsStreamAsync())
        using (StreamReader readStream = new StreamReader(stream, Encoding.UTF8))
        using (JsonTextReader jsonTextReader = new JsonTextReader(readStream))
        {
            JsonSerializer serializer = new JsonSerializer();
            ObjectDetails objectDetails = serializer.Deserialize<ObjectDetails>(
                jsonTextReader);
            return objectDetails;
        }
    }
    catch (Exception e)
    {
        // Log exception
        return null;
    }
}

private async Task<HttpResponseMessage> TryGetResponse(string urlStr)
{
    try
    {
        HttpResponseMessage response = await httpClient.GetAsync(urlStr)
            .ConfigureAwait(false);
        if (response.StatusCode != HttpStatusCode.OK)
        {
            throw new WebException("Response code is "
                + response.StatusCode.ToString() + "... not 200 OK.");
        }
        return response;
    }
    catch (Exception e)
    {
        // Log exception
        return null;
    }
}

private async Task<ListOfObjects> GetObjectDetailsAsync(string baseUrl, int id)
{
    string urlStr = baseUrl + @"objects/id/" + id + "/details";

    HttpResponseMessage response = await TryGetResponse(urlStr);

    ObjectDetails objectDetails = await TryDeserializeResponse(response);

    return objectDetails;
}

// With ~3000 objects to retrieve, this code will create 100 API calls
// in parallel, wait for all 100 to finish, and then repeat that process
// ~30 times. In other words, there will be ~30 batches of 100 parallel
// API calls.
private Dictionary<int, Task<ObjectDetails>> GetAllObjectDetailsInBatches(
    string baseUrl, Dictionary<int, MyObject> incompleteObjects)
{
    int batchSize = 100;
    int numberOfBatches = (int)Math.Ceiling(
        (double)incompleteObjects.Count / batchSize);
    Dictionary<int, Task<ObjectDetails>> objectTaskDict
        = new Dictionary<int, Task<ObjectDetails>>(incompleteObjects.Count);

    var orderedIncompleteObjects = incompleteObjects.OrderBy(pair => pair.Key);

    for (int i = 0; i < 1; i++)
    {
        var batchOfObjects = orderedIncompleteObjects.Skip(i * batchSize)
            .Take(batchSize);
        var batchObjectsTaskList = batchOfObjects.Select(
            pair => GetObjectDetailsAsync(baseUrl, pair.Key));
        Task.WaitAll(batchObjectsTaskList.ToArray());
        foreach (var objTask in batchObjectsTaskList)
            objectTaskDict.Add(objTask.Result.id, objTask);
    }

    return objectTaskDict;
}

public void GetObjectsVersion1()
{
    string baseUrl = @"https://mywebserver.com:/api";

    // GetIncompleteObjects is not shown, but it is not relevant to
    // the question
    Dictionary<int, MyObject> incompleteObjects = GetIncompleteObjects();

    Dictionary<int, Task<ObjectDetails>> objectTaskDict
        = GetAllObjectDetailsInBatches(baseUrl, incompleteObjects);

    foreach (KeyValuePair<int, MyObject> pair in incompleteObjects)
    {
        ObjectDetails objectDetails = objectTaskDict[pair.Key].Result
            .objectDetails;

        // Code here that copies fields from objectDetails to pair.Value
        // (the incompleteObject)

        AllObjects.Add(pair.Value);
    };
}

public void GetObjectsVersion2()
{
    string baseUrl = @"https://mywebserver.com:/api";

    // GetIncompleteObjects is not shown, but it is not relevant to
    // the question
    Dictionary<int, MyObject> incompleteObjects = GetIncompleteObjects();

    Parallel.ForEach(incompleteHosts, pair =>
    {
        ObjectDetails objectDetails = GetObjectDetailsAsync(
            baseUrl, pair.Key).Result.objectDetails;

        // Code here that copies fields from objectDetails to pair.Value
        // (the incompleteObject)

        AllObjects.Add(pair.Value);
    });
}

【问题讨论】:

  • 在某些时候您没有使用 ConfigureAwait(false)(请参阅 GetObjectDetailsAsync),这将对性能产生很大影响,因为代码正在等待同步。
  • 此外,这段代码在哪个上下文中运行也会很有趣。 WinForms/WPF/asp.net/Console/...?
  • @SirRufo 啊,是的,我将在那里添加 ConfigureAwait(false) 以查看它如何影响性能。它是一个控制台应用程序(.Net Framework 4.6.1)
  • @SirRufo 在 GetObjectDetailsAsync 中添加 ConfigureAwait(false) 似乎不会影响性能。
  • 是的,控制台应用程序根本没有同步上下文,因此没有会​​影响性能的同步。这就是我想知道代码在哪种应用程序中运行的原因

标签: c# asynchronous async-await task parallel.foreach


【解决方案1】:

Parallel.ForEach 可能运行得更快的一个可能原因是它会产生节流的副作用。最初 x 个线程正在处理前 x 个元素(其中 x 是可用内核的数量),并且可以根据内部启发式逐渐添加更多线程。限制 IO 操作是一件好事,因为它可以保护网络和处理请求的服务器不会变得负担过重。您的替代即兴节流方法,通过批量发出 100 个请求,由于许多原因远非理想,其中一个原因是 100 个并发请求是很多请求!另一个是单个长时间运行的操作可能会延迟批处理的完成,直到其他 99 个操作完成之后很久。

请注意,Parallel.ForEach 也不适合并行化 IO 操作。它只是碰巧比替代方案表现更好,一直在浪费内存。更好的方法请看这里:How to limit the amount of concurrent async I/O operations?

【讨论】:

  • 感谢您的精彩回答!因此,批处理解决方案与使用 SemaphoreSlim 的解决方案之间的区别在于,SemaphoreSlim 解决方案将始终同时运行 20 个请求 - 一旦 20 个请求中的一个完成就启动一个新请求,而批处理方法将等待启动下一批,直到当前批中的所有请求都完成。我喜欢它。
  • @davekats 是的,你有这个概念。 :-)
【解决方案2】:

https://docs.microsoft.com/en-us/dotnet/api/system.threading.tasks.parallel.foreach?view=netframework-4.8

基本上,并行 foreach 允许迭代并行运行,因此您不会将迭代限制为串行运行,在不受线程限制的主机上,这往往会提高吞吐量

【讨论】:

    【解决方案3】:

    简而言之:

    • Parallel.Foreach() 对于 CPU 密集型任务最有用。
    • Task.WaitAll() 对于 IO 绑定任务更有用。

    因此,在您的情况下,您从网络服务器(即 IO)获取信息。如果异步方法被正确实现,它不会阻塞任何线程。 (它将使用 IO 完成端口等待) 这样线程就可以做其他事情了。

    通过运行异步方法GetObjectDetailsAsync(baseUrl, pair.Key).Result同步,它会阻塞一个线程。所以线程池会被等待线程淹没。

    所以我认为 Task 解决方案会更合适。

    【讨论】:

    • 谢谢 Jeroen,我同意你的观点,但我的问题不是哪种解决方案更合适。我知道 Task 解决方案更合适,但我看到的是 Task 解决方案实际上更慢 - 一个意想不到的结果。所以我的问题是,为什么在这种情况下,并行解决方案更快?
    猜你喜欢
    • 2020-09-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-07-05
    • 2015-05-30
    • 1970-01-01
    • 2017-03-26
    • 2012-01-13
    相关资源
    最近更新 更多