.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。有关从异步方法传播异常的更多信息,请查看 here 或 here。
两种实现都能优雅地处理枚举source 期间可能发生的任何错误。 ForEachAsync 方法在所有挂起的操作完成之前不会完成。没有任何任务被遗漏(以即发即弃的方式)。
¹ 第一个实现elides async and await.