【问题标题】:Nesting await in Parallel.ForEach [duplicate]在 Parallel.ForEach 中嵌套 await [重复]
【发布时间】:2012-07-18 20:29:48
【问题描述】:

在 Metro 应用程序中,我需要执行多个 WCF 调用。需要进行大量调用,因此我需要在并行循环中进行调用。问题是并行循环在 WCF 调用全部完成之前退出。

您将如何重构它以按预期工作?

var ids = new List<string>() { "1", "2", "3", "4", "5", "6", "7", "8", "9", "10" };
var customers = new  System.Collections.Concurrent.BlockingCollection<Customer>();

Parallel.ForEach(ids, async i =>
{
    ICustomerRepo repo = new CustomerRepo();
    var cust = await repo.GetCustomer(i);
    customers.Add(cust);
});

foreach ( var customer in customers )
{
    Console.WriteLine(customer.ID);
}

Console.ReadKey();

【问题讨论】:

  • 我已将此问题投票为 Parallel foreach with asynchronous lambda 的重复项,尽管该问题比此问题更新了几个月,因为另一个问题包含一个已被高度支持的 answer 建议对于这个问题,目前最好的解决方案可能是什么,那就是新的Parallel.ForEachAsync API。

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


【解决方案1】:

Parallel.ForEach() 背后的整个想法是你有一组线程,每个线程处理集合的一部分。正如您所注意到的,这不适用于async-await,您希望在异步调用期间释放线程。

您可以通过阻塞ForEach() 线程来“解决”这个问题,但这违背了async-await 的全部意义。

你可以做的是使用TPL Dataflow而不是Parallel.ForEach(),它很好地支持异步Tasks。

具体来说,您的代码可以使用TransformBlock 编写,使用async lambda 将每个id 转换为Customer。该块可以配置为并行执行。您可以将该块链接到一个ActionBlock,该Customer 将每个Customer 写入控制台。 设置好区块网络后,可以将Post()每个id转为TransformBlock

在代码中:

var ids = new List<string> { "1", "2", "3", "4", "5", "6", "7", "8", "9", "10" };

var getCustomerBlock = new TransformBlock<string, Customer>(
    async i =>
    {
        ICustomerRepo repo = new CustomerRepo();
        return await repo.GetCustomer(i);
    }, new ExecutionDataflowBlockOptions
    {
        MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded
    });
var writeCustomerBlock = new ActionBlock<Customer>(c => Console.WriteLine(c.ID));
getCustomerBlock.LinkTo(
    writeCustomerBlock, new DataflowLinkOptions
    {
        PropagateCompletion = true
    });

foreach (var id in ids)
    getCustomerBlock.Post(id);

getCustomerBlock.Complete();
writeCustomerBlock.Completion.Wait();

尽管您可能希望将TransformBlock 的并行性限制为某个小常数。此外,您可以限制TransformBlock 的容量并使用SendAsync() 将项目异步添加到其中,例如如果集合太大。

与您的代码(如果有效)相比,另一个好处是,一旦完成单个项目,就会开始写入,而不是等到所有处理完成。

【讨论】:

  • 对异步、响应式扩展、TPL 和 TPL 数据流的简要概述 - vantsuyoshi.wordpress.com/2012/01/05/… 适合像我这样可能需要了解的人。
  • 我很确定这个答案不会并行处理。我相信您需要对 id 执行 Parallel.ForEach 并将其发布到 getCustomerBlock。至少这是我在测试这个建议时发现的。
  • @JasonLind 确实如此。并行使用 Parallel.ForEach()Post() 项目应该没有任何实际效果。
  • @svick 好的,我找到了,ActionBlock 也需要并行。我的做法略有不同,我不需要转换,所以我只使用了一个缓冲块并在 ActionBlock 中完成了我的工作。我对互联网上的另一个答案感到困惑。
  • 我的意思是在 ActionBlock 上指定 MaxDegreeOfParallelism,就像您在示例中的 TransformBlock 上所做的那样
【解决方案2】:

svick's answer(和往常一样)非常出色。

但是,当您实际上需要传输大量数据时,我发现 Dataflow 会更有用。或者当您需要async 兼容队列时。

在您的情况下,一个更简单的解决方案是只使用async-style 并行:

var ids = new List<string>() { "1", "2", "3", "4", "5", "6", "7", "8", "9", "10" };

var customerTasks = ids.Select(i =>
  {
    ICustomerRepo repo = new CustomerRepo();
    return repo.GetCustomer(i);
  });
var customers = await Task.WhenAll(customerTasks);

foreach (var customer in customers)
{
  Console.WriteLine(customer.ID);
}

Console.ReadKey();

【讨论】:

  • 如果你想手动限制并行度(在这种情况下你很可能会这样做),这样做会更复杂。
  • 但你说得对,Dataflow 可能相当复杂(例如与Parallel.ForEach() 相比)。但我认为这是目前几乎所有 async 处理集合的最佳选择。
  • @batmaci: Parallel.ForEach 不支持async
  • @MikeT:这不会按预期工作。 PLINQ 不理解异步任务,因此代码只会并行化 async lambda 的开始
  • @Mike: Parallel(和Task&lt;T&gt;)早于async/await 编写,作为任务并行库(TPL)的一部分。当async/await 出现时,他们可以选择制作自己的Future&lt;T&gt; 类型以与async 一起使用,或者重新使用TPL 中现有的Task&lt;T&gt; 类型。这两个决定显然都不正确,因此他们决定重新使用Task&lt;T&gt;
【解决方案3】:

按照 svick 的建议使用 DataFlow 可能有点矫枉过正,斯蒂芬的回答没有提供控制操作并发性的方法。然而,这可以很简单地实现:

public static async Task RunWithMaxDegreeOfConcurrency<T>(
     int maxDegreeOfConcurrency, IEnumerable<T> collection, Func<T, Task> taskFactory)
{
    var activeTasks = new List<Task>(maxDegreeOfConcurrency);
    foreach (var task in collection.Select(taskFactory))
    {
        activeTasks.Add(task);
        if (activeTasks.Count == maxDegreeOfConcurrency)
        {
            await Task.WhenAny(activeTasks.ToArray());
            //observe exceptions here
            activeTasks.RemoveAll(t => t.IsCompleted); 
        }
    }
    await Task.WhenAll(activeTasks.ToArray()).ContinueWith(t => 
    {
        //observe exceptions in a manner consistent with the above   
    });
}

ToArray() 调用可以通过使用数组而不是列表和替换已完成的任务来优化,但我怀疑它在大多数情况下会产生很大的不同。每个 OP 问题的示例用法:

RunWithMaxDegreeOfConcurrency(10, ids, async i =>
{
    ICustomerRepo repo = new CustomerRepo();
    var cust = await repo.GetCustomer(i);
    customers.Add(cust);
});

编辑 SO 用户和 TPL 向导 Eli Arbel 向我指出了 related article from Stephen Toub。像往常一样,他的实现既优雅又高效:

public static Task ForEachAsync<T>(
      this IEnumerable<T> source, int dop, Func<T, Task> body) 
{ 
    return Task.WhenAll( 
        from partition in Partitioner.Create(source).GetPartitions(dop) 
        select Task.Run(async delegate { 
            using (partition) 
                while (partition.MoveNext()) 
                    await body(partition.Current).ContinueWith(t => 
                          {
                              //observe exceptions
                          });
                      
        })); 
}

【讨论】:

  • @RichardPierre 实际上这种Partitioner.Create 的重载使用了块分区,它为不同的任务动态地提供元素,因此您描述的场景不会发生。另请注意,由于开销较小(特别是同步),静态(预定)分区在某些情况下可能会更快。欲了解更多信息,请参阅:msdn.microsoft.com/en-us/library/dd997411(v=vs.110).aspx
  • @OhadSchneider 在 // 观察异常中,如果抛出异常,它会冒泡到调用者吗?例如,如果我希望整个可枚举在它的任何部分失败时停止处理/失败?
  • @Terry 它将冒泡到调用者,因为最顶层的任务(由Task.WhenAll 创建)将包含异常(在AggregateException 内),因此如果说调用者使用await,调用站点会抛出异常。但是,Task.WhenAll 仍将等待所有 任务完成,并且GetPartitions 将在调用partition.MoveNext 时动态分配元素,直到没有更多元素需要处理。这意味着除非您添加自己的机制来停止处理(例如CancellationToken),否则它不会自行发生。
  • @MichaelFreidgeim 您可以在await body 之前执行var current = partition.Current 之类的操作,然后在继续使用current (ContinueWith(t =&gt; { ... })。
  • Stephen Toub 文章的更新链接:devblogs.microsoft.com/pfxteam/…
【解决方案4】:

您可以使用新的AsyncEnumerator NuGet Package 来节省精力,4 年前最初发布问题时它还不存在。它允许您控制并行度:

using System.Collections.Async;
...

await ids.ParallelForEachAsync(async i =>
{
    ICustomerRepo repo = new CustomerRepo();
    var cust = await repo.GetCustomer(i);
    customers.Add(cust);
},
maxDegreeOfParallelism: 10);

免责声明:我是 AsyncEnumerator 库的作者,该库是开源的并在 MIT 下获得许可,我发布此消息只是为了帮助社区。​​p>

【讨论】:

  • Sergey,你应该公开你是图书馆的作者
  • 好的,添加了免责声明。我不是从广告中寻求任何好处,只是想帮助人们;)
  • 您的库与 .NET Core 不兼容。
  • @CornielNobel,它与 .NET Core 兼容 - GitHub 上的源代码具有 .NET Framework 和 .NET Core 的测试覆盖率。
  • @SergeSemenov 我经常使用你的库,因为它的AsyncStreams 我不得不说它非常棒。不能推荐这个库。
【解决方案5】:

Parallel.Foreach 包装成Task.Run(),而不是await 关键字使用[yourasyncmethod].Result

(您需要执行 Task.Run 操作以不阻塞 UI 线程)

类似这样的:

var yourForeachTask = Task.Run(() =>
        {
            Parallel.ForEach(ids, i =>
            {
                ICustomerRepo repo = new CustomerRepo();
                var cust = repo.GetCustomer(i).Result;
                customers.Add(cust);
            });
        });
await yourForeachTask;

【讨论】:

  • 这有什么问题?我会这样做的。让Parallel.ForEach 做并行工作,这会阻塞直到所有工作完成,然后将整个事情推送到后台线程以获得响应式 UI。有什么问题吗?也许这是一个睡眠线程太多了,但它是简短易读的代码。
  • @LonelyPixel 我唯一的问题是,当TaskCompletionSource 更可取时,它会调用Task.Run
  • @Gusdor 好奇 - 为什么TaskCompletionSource 更可取?
  • 只是一个简短的更新。我现在正在寻找这个,向下滚动以找到最简单的解决方案并再次找到我自己的评论。我完全使用了这段代码,它按预期工作。它只假设循环中有原始异步调用的同步版本。 await可以移到前面来保存多余的变量名。
  • 我不确定您的方案是什么,但我相信您可以删除 Task.Run()。只需在末尾附加一个 .Result 或 .Wait 就足以使并行执行等待所有线程完成。
【解决方案6】:

这应该非常有效,并且比让整个 TPL 数据流正常工作更容易:

var customers = await ids.SelectAsync(async i =>
{
    ICustomerRepo repo = new CustomerRepo();
    return await repo.GetCustomer(i);
});

...

public static async Task<IList<TResult>> SelectAsync<TSource, TResult>(this IEnumerable<TSource> source, Func<TSource, Task<TResult>> selector, int maxDegreesOfParallelism = 4)
{
    var results = new List<TResult>();

    var activeTasks = new HashSet<Task<TResult>>();
    foreach (var item in source)
    {
        activeTasks.Add(selector(item));
        if (activeTasks.Count >= maxDegreesOfParallelism)
        {
            var completed = await Task.WhenAny(activeTasks);
            activeTasks.Remove(completed);
            results.Add(completed.Result);
        }
    }

    results.AddRange(await Task.WhenAll(activeTasks));
    return results;
}

【讨论】:

  • 该用法示例不应该使用await like:var customers = await ids.SelectAsync(async i =&gt; { ... });吗?
【解决方案7】:

一种利用 SemaphoreSlim 并允许设置最大并行度的扩展方法

    /// <summary>
    /// Concurrently Executes async actions for each item of <see cref="IEnumerable<typeparamref name="T"/>
    /// </summary>
    /// <typeparam name="T">Type of IEnumerable</typeparam>
    /// <param name="enumerable">instance of <see cref="IEnumerable<typeparamref name="T"/>"/></param>
    /// <param name="action">an async <see cref="Action" /> to execute</param>
    /// <param name="maxDegreeOfParallelism">Optional, An integer that represents the maximum degree of parallelism,
    /// Must be grater than 0</param>
    /// <returns>A Task representing an async operation</returns>
    /// <exception cref="ArgumentOutOfRangeException">If the maxActionsToRunInParallel is less than 1</exception>
    public static async Task ForEachAsyncConcurrent<T>(
        this IEnumerable<T> enumerable,
        Func<T, Task> action,
        int? maxDegreeOfParallelism = null)
    {
        if (maxDegreeOfParallelism.HasValue)
        {
            using (var semaphoreSlim = new SemaphoreSlim(
                maxDegreeOfParallelism.Value, maxDegreeOfParallelism.Value))
            {
                var tasksWithThrottler = new List<Task>();

                foreach (var item in enumerable)
                {
                    // Increment the number of currently running tasks and wait if they are more than limit.
                    await semaphoreSlim.WaitAsync();

                    tasksWithThrottler.Add(Task.Run(async () =>
                    {
                        await action(item).ContinueWith(res =>
                        {
                            // action is completed, so decrement the number of currently running tasks
                            semaphoreSlim.Release();
                        });
                    }));
                }

                // Wait for all tasks to complete.
                await Task.WhenAll(tasksWithThrottler.ToArray());
            }
        }
        else
        {
            await Task.WhenAll(enumerable.Select(item => action(item)));
        }
    }

示例用法:

await enumerable.ForEachAsyncConcurrent(
    async item =>
    {
        await SomeAsyncMethod(item);
    },
    5);

【讨论】:

    【解决方案8】:

    我参加聚会有点晚了,但您可能想考虑使用 GetAwaiter.GetResult() 在同步上下文中运行您的异步代码,但并行如下;

     Parallel.ForEach(ids, i =>
    {
        ICustomerRepo repo = new CustomerRepo();
        // Run this in thread which Parallel library occupied.
        var cust = repo.GetCustomer(i).GetAwaiter().GetResult();
        customers.Add(cust);
    });
    

    【讨论】:

      【解决方案9】:

      在引入一堆辅助方法后,您将能够使用以下简单语法运行并行查询:

      const int DegreeOfParallelism = 10;
      IEnumerable<double> result = await Enumerable.Range(0, 1000000)
          .Split(DegreeOfParallelism)
          .SelectManyAsync(async i => await CalculateAsync(i).ConfigureAwait(false))
          .ConfigureAwait(false);
      

      这里发生的情况是:我们将源集合拆分为 10 个块 (.Split(DegreeOfParallelism)),然后运行 ​​10 个任务,每个任务一个接一个地处理其项目 (.SelectManyAsync(...)),并将它们合并回一个列表。

      值得一提的是有一个更简单的方法:

      double[] result2 = await Enumerable.Range(0, 1000000)
          .Select(async i => await CalculateAsync(i).ConfigureAwait(false))
          .WhenAll()
          .ConfigureAwait(false);
      

      但需要注意事项:如果您的源集合太大,它会立即为每个项目安排Task,这可能会导致严重的性能下降。

      以上示例中使用的扩展方法如下所示:

      public static class CollectionExtensions
      {
          /// <summary>
          /// Splits collection into number of collections of nearly equal size.
          /// </summary>
          public static IEnumerable<List<T>> Split<T>(this IEnumerable<T> src, int slicesCount)
          {
              if (slicesCount <= 0) throw new ArgumentOutOfRangeException(nameof(slicesCount));
      
              List<T> source = src.ToList();
              var sourceIndex = 0;
              for (var targetIndex = 0; targetIndex < slicesCount; targetIndex++)
              {
                  var list = new List<T>();
                  int itemsLeft = source.Count - targetIndex;
                  while (slicesCount * list.Count < itemsLeft)
                  {
                      list.Add(source[sourceIndex++]);
                  }
      
                  yield return list;
              }
          }
      
          /// <summary>
          /// Takes collection of collections, projects those in parallel and merges results.
          /// </summary>
          public static async Task<IEnumerable<TResult>> SelectManyAsync<T, TResult>(
              this IEnumerable<IEnumerable<T>> source,
              Func<T, Task<TResult>> func)
          {
              List<TResult>[] slices = await source
                  .Select(async slice => await slice.SelectListAsync(func).ConfigureAwait(false))
                  .WhenAll()
                  .ConfigureAwait(false);
              return slices.SelectMany(s => s);
          }
      
          /// <summary>Runs selector and awaits results.</summary>
          public static async Task<List<TResult>> SelectListAsync<TSource, TResult>(this IEnumerable<TSource> source, Func<TSource, Task<TResult>> selector)
          {
              List<TResult> result = new List<TResult>();
              foreach (TSource source1 in source)
              {
                  TResult result1 = await selector(source1).ConfigureAwait(false);
                  result.Add(result1);
              }
              return result;
          }
      
          /// <summary>Wraps tasks with Task.WhenAll.</summary>
          public static Task<TResult[]> WhenAll<TResult>(this IEnumerable<Task<TResult>> source)
          {
              return Task.WhenAll<TResult>(source);
          }
      }
      

      【讨论】:

        【解决方案10】:

        .NET 6 更新: 引入Parallel.ForEachAsync API 后,以下实现不再相关。它们仅对面向 .NET 6 之前的 .NET 平台版本的项目有用。


        下面是 ForEachAsync 方法的简单通用实现,它基于来自 TPL Dataflow 库的 ActionBlock,现在嵌入在 .NET 5 平台中:

        public static Task ForEachAsync<T>(this IEnumerable<T> source,
            Func<T, Task> action, int dop)
        {
            // Arguments validation omitted
            var block = new ActionBlock<T>(action,
                new ExecutionDataflowBlockOptions() { MaxDegreeOfParallelism = dop });
            try
            {
                foreach (var item in source) block.Post(item);
                block.Complete();
            }
            catch (Exception ex) { ((IDataflowBlock)block).Fault(ex); }
            return block.Completion;
        }
        

        此解决方案急切地枚举提供的IEnumerable,并立即将其所有元素发送到ActionBlock。所以它不太适合具有大量元素的可枚举。下面是一种更复杂的方法,它懒惰地枚举源,并将其元素一个一个发送到ActionBlock

        public static async Task ForEachAsync<T>(this IEnumerable<T> source,
            Func<T, Task> action, int dop)
        {
            // Arguments validation omitted
            var block = new ActionBlock<T>(action, new ExecutionDataflowBlockOptions()
            { MaxDegreeOfParallelism = dop, BoundedCapacity = dop });
            try
            {
                foreach (var item in source)
                    if (!await block.SendAsync(item).ConfigureAwait(false)) break;
                block.Complete();
            }
            catch (Exception ex) { ((IDataflowBlock)block).Fault(ex); }
            try { await block.Completion.ConfigureAwait(false); }
            catch { block.Completion.Wait(); } // Propagate AggregateException
        }
        

        这两种方法在出现异常时具有不同的行为。第一个¹ 在其 InnerExceptions 属性中直接传播包含异常的 AggregateException。第二个传播一个AggregateException,其中包含另一个AggregateException,但有例外。就我个人而言,我发现第二种方法的行为在实践中更方便,因为等待它会自动消除一层嵌套,所以我可以简单地catch (AggregateException aex) 并在catch 块内处理aex.InnerExceptions。第一种方法需要在等待之前存储Task,以便我可以访问catch 块内的task.Exception.InnerExceptions。有关从异步方法传播异常的更多信息,请查看 herehere

        两种实现都能优雅地处理枚举source 期间可能发生的任何错误。 ForEachAsync 方法在所有挂起的操作完成之前不会完成。没有任何任务被遗漏(以即发即弃的方式)。

        ¹ 第一个实现elides async and await.

        【讨论】:

        • 这与您共享的其他ForEachAsync() 实现here 相比如何?
        • @alhazen 这个实现在功能上与the other implementation 相同,假设默认行为bool onErrorContinue = false。此实现利用了 TPL Dataflow 库,因此代码更短,包含未发现错误的概率更小。在性能方面,这两个实现也应该非常相似。
        • @alhazen 实际上是有区别的。此实现在 ThreadPool 上调用异步委托 (Func&lt;T, Task&gt; action),而 the other implementation 在当前上下文中调用它。因此,例如,如果委托访问 UI 组件(假设是 WPF/WinForms 应用程序),则此实现很可能会失败,而另一个将按预期工作。
        【解决方案11】:

        没有 TPL 的简单原生方式:

        int totalThreads = 0; int maxThreads = 3;
        
        foreach (var item in YouList)
        {
            while (totalThreads >= maxThreads) await Task.Delay(500);
            Interlocked.Increment(ref totalThreads);
        
            MyAsyncTask(item).ContinueWith((res) => Interlocked.Decrement(ref totalThreads));
        }
        

        您可以在下一个任务中检查此解决方案:

        async static Task MyAsyncTask(string item)
        {
            await Task.Delay(2500);
            Console.WriteLine(item);
        }
        

        【讨论】:

        • 不错的尝试,但这种方法存在多个问题:在不同步的情况下访问非volatile 变量totalThreads。在循环中无效率地等待满足条件(引入延迟)。使用 primitive ContinueWith 方法而不指定 TaskScheduler。在MyAsyncTask 同步抛出的情况下,泄漏即发即弃任务的可能性。这个功能出奇的棘手,而且你不可能在第一次尝试时就自己做。
        猜你喜欢
        • 1970-01-01
        • 2012-09-19
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-04-04
        • 2015-11-10
        • 2015-09-09
        相关资源
        最近更新 更多