如何让Parallel循环等待内部异步Process方法执行完成?
问题:Parallel.ForEach未等待异步Process方法完成的原因及解决办法
代码示例
var index = 0; var options = new ParallelOptions { MaxDegreeOfParallelism = maxParallelCount }; ParallelLoopResult result = Parallel.ForEach(customers, options, async cust => { var currentCount = Interlocked.Increment(ref index); Console.WriteLine("Processing {0}/{1}", currentCount, totalCount); var success = await Process(cust); if (success) Console.WriteLine("({0}/{1}) SUCCESS", cust.Name, cust.Number); else Console.WriteLine("({0}/{1}) FAILED", cust.Name, cust.Number); }); if (result.IsCompleted) Console.WriteLine("process completed | {0} customers processed.", totalCount);
实际输出
process completed | 10 customers processed. Customer Name1/Number1 SUCCESS Customer Name1/Number4 SUCCESS Customer Name1/Number3 SUCCESS Customer Name1/Number2 SUCCESS etc..
预期输出
Customer Name1/Number1 SUCCESS Customer Name1/Number4 SUCCESS Customer Name1/Number3 SUCCESS Customer Name1/Number2 SUCCESS etc.. process completed | 10 customers processed.
疑问
为何ParallelLoopResult未等待Process()方法执行完成?如何改进代码使其等待Process()执行完成?
原因分析
Parallel.ForEach的委托参数是Action<T>,当你传入异步lambda表达式时,它会被隐式转换为async void方法。async void属于异步编程里的“火即忘”模式,Parallel.ForEach只会等待lambda中await之前的同步代码执行完毕,就判定该循环迭代完成,不会等待后续异步的Process方法执行。因此result.IsCompleted会提前变为true,导致“完成”日志先输出,而Process的结果日志要等后台异步操作结束后才陆续打印。
改进方案
方案1:使用Task.WhenAll结合并发控制(兼容所有.NET版本)
把每个客户的处理逻辑包装成Task,用Task.WhenAll等待所有任务完成,同时通过SemaphoreSlim限制并发数,替代原代码中ParallelOptions.MaxDegreeOfParallelism的作用:
var index = 0; var semaphore = new SemaphoreSlim(maxParallelCount); var tasks = new List<Task>(); foreach (var cust in customers) { await semaphore.WaitAsync(); tasks.Add(Task.Run(async () => { try { var currentCount = Interlocked.Increment(ref index); Console.WriteLine("Processing {0}/{1}", currentCount, totalCount); var success = await Process(cust); if (success) Console.WriteLine("({0}/{1}) SUCCESS", cust.Name, cust.Number); else Console.WriteLine("({0}/{1}) FAILED", cust.Name, cust.Number); } finally { semaphore.Release(); } })); } await Task.WhenAll(tasks); Console.WriteLine("process completed | {0} customers processed.", totalCount);
方案2:使用Parallel.ForEachAsync(.NET 6及以上版本推荐)
.NET 6新增了原生支持异步的Parallel.ForEachAsync方法,它可以直接处理异步委托,同时支持并发度控制,用法和原代码接近:
var index = 0; var options = new ParallelOptions { MaxDegreeOfParallelism = maxParallelCount }; await Parallel.ForEachAsync(customers, options, async (cust, cancellationToken) => { var currentCount = Interlocked.Increment(ref index); Console.WriteLine("Processing {0}/{1}", currentCount, totalCount); var success = await Process(cust); if (success) Console.WriteLine("({0}/{1}) SUCCESS", cust.Name, cust.Number); else Console.WriteLine("({0}/{1}) FAILED", cust.Name, cust.Number); }); Console.WriteLine("process completed | {0} customers processed.", totalCount);
内容的提问来源于stack exchange,提问作者Jyina
相关产品推荐
相关产品推荐

