使用Async Await处理大量LongRunning并行任务时出现挂起问题
问题描述
尝试同时向REST服务发起500个请求,编写了一个扩展方法,该方法有时正常运行,但有时会在Task.WhenAll处挂起,存在未完成的任务。去掉LongRunning标记后,每次都会挂起(或耗时远超预期),推测是线程池过载导致。正常运行时该方法仅需约4秒。
原实现代码:
public static async Task<List<TFunc>> ForEachLongRunningAsync<TSource, TFunc>(this IEnumerable<TSource> source, Func<TSource, Task<TFunc>> func) { var tasks = source.Select(x => Task.Factory.StartNew(() => func(x), TaskCreationOptions.LongRunning | TaskCreationOptions.AttachedToParent).Unwrap()); var results = await Task.WhenAll(tasks); return results.ToList(); }
为排查问题修改后的代码,仍时而正常时而挂起:
public static async Task<List<TFunc>> ForEachLongRunningAsync<TSource, TFunc>(this IEnumerable<TSource> source, Func<TSource, Task<TFunc>> func) { var tasks = source.Select(x => Task.Factory.StartNew(() => func(x), TaskCreationOptions.LongRunning | TaskCreationOptions.AttachedToParent).Unwrap()).ToList(); var results = Task.WhenAll(tasks); while (true) { if (results.IsCompleted) break; await Task.Delay(100); } return results.Result.ToList(); }
原因分析
LongRunning与AttachedToParent的冲突:AttachedToParent会将子任务附加到父任务,但LongRunning创建的是独立线程的任务,这种组合可能导致任务状态追踪异常,部分任务的完成信号无法正确传递到Task.WhenAll。- 不必要的线程创建:REST请求是IO密集型操作,用
LongRunning创建额外线程完全没必要,反而会浪费系统资源,甚至因为线程过多导致调度问题。 - 线程池过载的误解:去掉
LongRunning后挂起,不是因为线程池过载——IO密集型任务在等待响应时会释放线程池线程,真正的问题是原代码用Task.Factory.StartNew包裹异步方法,额外的线程调度反而引入了阻塞点。
最优实现方案
基础版(无并发限制)
直接遍历生成异步任务,不需要额外用Task.Factory.StartNew包装,因为func本身就是返回Task的异步方法,IO密集型任务天然适合无额外线程的异步调度:
public static async Task<List<TFunc>> ForEachAsync<TSource, TFunc>(this IEnumerable<TSource> source, Func<TSource, Task<TFunc>> func) { var tasks = source.Select(func).ToList(); var results = await Task.WhenAll(tasks); return results.ToList(); }
带并发限制版(推荐)
一次性发起500个请求可能压垮目标服务或本地网络,用SemaphoreSlim做流量控制,平衡并发效率与服务稳定性:
public static async Task<List<TFunc>> ForEachAsync<TSource, TFunc>(this IEnumerable<TSource> source, Func<TSource, Task<TFunc>> func, int maxConcurrency) { using var semaphore = new SemaphoreSlim(maxConcurrency); var tasks = source.Select(async item => { await semaphore.WaitAsync(); try { return await func(item); } finally { semaphore.Release(); } }).ToList(); var results = await Task.WhenAll(tasks); return results.ToList(); }
补充说明
- 对于IO密集型操作,异步编程的核心是避免阻塞线程,不需要额外创建线程。
Task.Factory.StartNew适合CPU密集型任务,用来包装异步方法只会增加不必要的开销。 - 限制并发数不仅能避免本地线程调度问题,还能遵守REST服务的限流规则,减少请求失败的概率。
内容的提问来源于stack exchange,提问作者LorneCash
相关产品推荐
相关产品推荐

