如何避免Task.WhenAll执行大量任务后出现超时?
问题描述
我有一段从Controller调用的代码如下:
public async Task Execute() { var collections= await _repo.GetCollections(); // 获取500+条数据 List<Object1> coolCollections= new List<Object1>(); List<Object2> uncoolCollections= new List<Object2>(); foreach (var collection in collections) { if(collection == "Something") { var specialObject = TurnObjectIntoSpecialObject(collection); uncoolCollections.Add(specialObject); } else { var anotherObject = TurnObjectIntoAnotherObject(collection); // 修正原代码笔误:原参数为object,应为collection coolCollections.Add(anotherObject); } } var list1Async = coolCollections.Select(async obj => await restService.PostObject1(obj)); // 每个请求耗时200-2000ms var list2Async = uncoolCollections.Select(async obj => await restService.PostObject2(obj));// 每个请求耗时300-3000ms var asyncTasks = list1Async.Concat<Task>(list2Async); await Task.WhenAll(asyncTasks); // 同时发起500+请求 }
执行约300个请求后遇到504错误,无法修改被调用的API,仅能优化上述代码。改用foreach串行执行可解决超时,但速度极慢。需要找到既能避免超时、又能保证执行效率的优化方案。
解决方案
核心问题是一次性发起的并发请求数量过多,超出了目标API的承载上限,导致网关或服务端返回504超时。需要通过控制并发请求数量来平衡执行效率和服务端压力,以下是两种实用方案:
方案1:用SemaphoreSlim做动态并发限流
通过信号量控制同时运行的请求数,既避免瞬间请求过载,又能最大化利用并发资源:
public async Task Execute() { var collections = await _repo.GetCollections(); // 设置并发上限,建议从20-50开始测试,根据实际情况调整 var concurrencySemaphore = new SemaphoreSlim(initialCount: 50); var tasks = new List<Task>(); foreach (var collection in collections) { // 等待获取信号量,超过并发上限时会阻塞 await concurrencySemaphore.WaitAsync(); tasks.Add(Task.Run(async () => { try { if (collection == "Something") { var specialObject = TurnObjectIntoSpecialObject(collection); await restService.PostObject2(specialObject); } else { var anotherObject = TurnObjectIntoAnotherObject(collection); await restService.PostObject1(anotherObject); } } finally { // 无论请求成功或失败,都释放信号量 concurrencySemaphore.Release(); } })); } await Task.WhenAll(tasks); }
方案2:分批次批量执行
将请求分成固定大小的批次,完成一批后再执行下一批,逻辑更直观:
public async Task Execute() { var collections = await _repo.GetCollections(); const int batchSize = 50; // 每批处理的请求数量,可按需调整 for (int i = 0; i < collections.Count; i += batchSize) { var currentBatch = collections.Skip(i).Take(batchSize); var batchTasks = currentBatch.Select(async collection => { if (collection == "Something") { var specialObject = TurnObjectIntoSpecialObject(collection); await restService.PostObject2(specialObject); } else { var anotherObject = TurnObjectIntoAnotherObject(collection); await restService.PostObject1(anotherObject); } }); // 等待当前批次所有请求完成,再执行下一批 await Task.WhenAll(batchTasks); } }
额外优化建议
- 添加重试机制:针对偶发的504超时,可引入重试逻辑(例如使用Polly库),避免单个请求失败影响整体流程:
// 定义重试策略:最多重试3次,每次间隔指数递增 var retryPolicy = Policy .Handle<HttpRequestException>() .OrResult<HttpResponseMessage>(r => r.StatusCode == HttpStatusCode.GatewayTimeout) .WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromSeconds(Math.Pow(2, retryAttempt))); // 在请求时应用策略 await retryPolicy.ExecuteAsync(() => restService.PostObject1(anotherObject)); - 调优并发数/批次大小:没有固定最优值,需根据目标API的实际承载能力测试调整——若仍出现504则降低数值,若速度过慢则适当提高。
- 减少中间集合:原代码中先将数据分成两个List再发起请求,可直接在遍历collections时发起请求,减少内存占用和不必要的中间步骤。
内容的提问来源于stack exchange,提问作者Greg Jones
相关产品推荐
相关产品推荐

