如何在Parallel.ForEach中等待异步函数执行完成?
异步并行任务等待问题的修复
问题重现
现有以下用于在字典中收集计算结果的函数:
void Func() { // ... var dtDict = await HandleComputeBooster(); // ... } async private static Task DoBooster(..., ConcurrentDictionary<string, DataTable> dtDict, ...) { DataTable dt = ...; // ... 数据处理逻辑 dtDict[symbol] = dt; // ... }
HandleComputeBooster函数存在提前返回的问题:
async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster() { var dtDict = new ConcurrentDictionary<string, DataTable>(); // ... 初始化逻辑 var chunks = listOfBoosterSymbols.ChunkBy(8); var pcCount = Environment.ProcessorCount; Parallel.ForEach(chunks, new ParallelOptions { MaxDegreeOfParallelism = pcCount - 2 }, async listStr => { var symbol = listStr[0]; await DoBooster(..., dtDict, ...); } ); // ... return dtDict; }
问题在于:HandleComputeBooster会在dtDict的所有计算任务完成前就返回。虽然最终所有数据都会存入字典,但需要让函数等待所有块处理完成后再返回。
问题原因
Parallel.ForEach不支持异步委托:传入的async lambda会被当作async void执行。Parallel.ForEach仅会等待线程池中的任务启动完成,不会等待委托内部的await操作结束,因此函数会提前返回,此时字典中的数据还未全部生成。
解决方案
改用Task.WhenAll结合异步遍历的方式,来正确等待所有异步任务完成:
async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster() { var dtDict = new ConcurrentDictionary<string, DataTable>(); // ... 初始化逻辑 var chunks = listOfBoosterSymbols.ChunkBy(8); // 将每个chunk的处理包装为Task var tasks = chunks.Select(async listStr => { var symbol = listStr[0]; await DoBooster(..., dtDict, ...); }); // 等待所有异步任务完成 await Task.WhenAll(tasks); // ... return dtDict; }
可选:限制并发数
如果需要像原代码一样限制最大并发数,可以使用SemaphoreSlim来控制:
async private Task<ConcurrentDictionary<string, DataTable>> HandleComputeBooster() { var dtDict = new ConcurrentDictionary<string, DataTable>(); // ... 初始化逻辑 var chunks = listOfBoosterSymbols.ChunkBy(8); var pcCount = Environment.ProcessorCount; var semaphore = new SemaphoreSlim(pcCount - 2); var tasks = chunks.Select(async listStr => { await semaphore.WaitAsync(); try { var symbol = listStr[0]; await DoBooster(..., dtDict, ...); } finally { semaphore.Release(); } }); await Task.WhenAll(tasks); semaphore.Dispose(); // ... return dtDict; }
说明
- 因为
DoBooster是异步方法,属于IO密集型操作,使用异步并行(Task.WhenAll)比Parallel.ForEach更合适,能更高效利用线程资源。 ConcurrentDictionary本身是线程安全的,多个异步任务同时写入不会有问题。
内容的提问来源于stack exchange,提问作者Ivan
相关产品推荐
相关产品推荐

