使用await仍出现Task内代码未完成就执行外部代码的问题解决
问题描述
编写了一个GetStatus方法,接收BulkStatusQueryMessage类型参数,通过Parallel.ForEach并行处理其中的StatusQueries集合,每个元素的处理逻辑中调用异步方法SendApiGetCall,并用await Task.Run包裹该并行操作。调试发现,即便使用了await,Task内部的代码(比如异常捕获逻辑)还没执行完,外部的return语句就先执行了——比如SendApiGetCall抛出异常时,会先进入return代码行,之后才触发catch逻辑。
相关代码如下:
private async Task<BulkResults> GetStatus(BulkStatusQueryMessage bulkStatusQueryMessage) { var results = new ConcurrentBag<Results>(); await Task.Run(async () => Parallel.ForEach(bulkStatusQueryMessage.StatusQueries, async sq => { try { var res = await _apiMethods.SendApiGetCall(ApiCallType.Status, queryParameters: sq.Id); var result = ExtractResultsFromResponse(res.Response, sq.OriginalRequestReferenceId); results.Add(result); } catch (Exception ex) { _logger.Error(ex, "Failed to get status for id {sq.id}" + " with OriginalRequestReferenceId {sq.OriginalRequestReferenceId}"); } })); return new BulkResults { ReferenceId = bulkStatusQueryMessage.ReferenceId, Responses = results }; }
注:BulkStatusQueryMessage类包含referenceId属性,以及一个包含id和OriginalRequestReferenceId属性的内部类集合。
问题原因分析
Parallel.ForEach不支持异步委托:Parallel.ForEach的委托参数如果是异步方法(async sq => ...),它会把这个异步委托当成async void来处理——这意味着Parallel.ForEach会直接启动每个异步任务,但不会等待它们完成。Task.Run的await仅等待Parallel.ForEach本身完成,而非内部异步任务:Task.Run包裹的是Parallel.ForEach的执行,Parallel.ForEach遍历完所有元素就会返回,此时内部的await _apiMethods.SendApiGetCall可能还没执行完,所以外部的await Task.Run会提前完成,进而执行后续的return语句,而内部的异步操作(包括异常捕获)还在后台运行。
修正方案
方案一:使用Task.WhenAll实现异步并行(推荐)
放弃Parallel.ForEach,改用Task.WhenAll来等待所有异步任务完成,这是异步场景下并行处理的标准做法,无需使用ConcurrentBag(因为每个任务返回结果后直接收集):
private async Task<BulkResults> GetStatus(BulkStatusQueryMessage bulkStatusQueryMessage) { // 为每个查询创建异步任务 var tasks = bulkStatusQueryMessage.StatusQueries.Select(async sq => { try { var res = await _apiMethods.SendApiGetCall(ApiCallType.Status, queryParameters: sq.Id); return ExtractResultsFromResponse(res.Response, sq.OriginalRequestReferenceId); } catch (Exception ex) { _logger.Error(ex, "Failed to get status for id {Id} with OriginalRequestReferenceId {OriginalRequestReferenceId}", sq.Id, sq.OriginalRequestReferenceId); // 返回null或者一个标记为失败的Results对象,根据业务需求调整 return null; } }); // 等待所有任务完成并收集结果,过滤掉null(如果有的话) var results = await Task.WhenAll(tasks); var validResults = results.Where(r => r != null).ToList(); return new BulkResults { ReferenceId = bulkStatusQueryMessage.ReferenceId, Responses = validResults }; }
方案二:如果必须使用Parallel.ForEach(不推荐异步场景)
如果坚持用Parallel.ForEach,需要手动等待每个异步任务完成,避免async void的问题,比如使用.GetAwaiter().GetResult()(但会阻塞线程,失去异步优势):
private async Task<BulkResults> GetStatus(BulkStatusQueryMessage bulkStatusQueryMessage) { var results = new ConcurrentBag<Results>(); await Task.Run(() => Parallel.ForEach(bulkStatusQueryMessage.StatusQueries, sq => { try { // 同步等待异步方法完成,注意:这会阻塞线程,不推荐 var res = _apiMethods.SendApiGetCall(ApiCallType.Status, queryParameters: sq.Id).GetAwaiter().GetResult(); var result = ExtractResultsFromResponse(res.Response, sq.OriginalRequestReferenceId); results.Add(result); } catch (Exception ex) { _logger.Error(ex, "Failed to get status for id {Id} with OriginalRequestReferenceId {OriginalRequestReferenceId}", sq.Id, sq.OriginalRequestReferenceId); } })); return new BulkResults { ReferenceId = bulkStatusQueryMessage.ReferenceId, Responses = results }; }
内容的提问来源于stack exchange,提问作者PMO1948
相关产品推荐
相关产品推荐

