Parallel.ForEachAsync与CancellationToken异常不一致问题排查
异步多线程与CancellationToken行为疑问
我对多线程async/await编程经验不足,代码中多处await的耗时不确定,导致日志和行为不一致,因此产生以下疑问。我的系统架构为:网站API端点 → 代理API → 外部API。当前核心困惑是理解CancellationToken的工作机制,以及为何OperationCanceledException的抛出/日志记录不符合预期。我编写了一个.NET 7控制台应用模拟该架构,代码如下。
我的问题
- 取消注释
throw new ApplicationException($"Site - Fake Exception for {item}");并反复运行程序时:- 从未收到代理API的日志,仿佛在
httpClient.SendAsync调用被取消前,请求从未到达代理API。为何准备HTTP请求耗时如此之长,导致CancellationToken已被取消? - 并非总能收到10条“Site -> Api Proxy”日志(即异步委托中的
logs.Enqueue($"Site: Calling ProxyApiAsync for {item}");行)。原因是什么?
- 从未收到代理API的日志,仿佛在
- 注释掉
throw new ApplicationException($"Site - Fake Exception for {item}");并反复运行程序时:- 成功收到10条“Site -> Api Proxy”和“Api Proxy -> External”日志,但并非总能在网站和代理API中收到9条“OperationCanceledException捕获”日志。如果向代理API提交10个请求且其中1个失败,难道不应该始终收到9条取消日志吗?(我猜测缺失的日志代表那些在任务3抛出异常前已成功完成的代理→外部请求?)
模拟代码
using System; using System.Linq; using System.Threading; using Microsoft.AspNetCore.Builder; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using System.Net.Http; using System.Threading.Tasks; using System.Collections.Concurrent; public static class Program { // Use this instead of Console.WriteLine() to enable some 'summary' queries after processing static ConcurrentQueue<string> logs = new ConcurrentQueue<string>(); public static async Task Main() { logs = new ConcurrentQueue<string>(); var builder = WebApplication.CreateBuilder(); builder .Logging.SetMinimumLevel(LogLevel.None) .Services.AddHttpClient(); var app = builder.Build(); app.MapGet("/proxyapi/{taskId}", ProxyApiAsync); // simulated Proxy Api web api var minApiTask = app.RunAsync(); await Task.Delay(100); // Wait for the API to start var httpClientFactory = app.Services.GetRequiredService<IHttpClientFactory>(); // simulated support for DI var apiInfosToRun = Enumerable.Range(0, 10).ToArray(); // simulated list of APIs to call try { await WebSiteEndpointAsync( apiInfosToRun, httpClientFactory ); } catch (Exception ex) { logs.Enqueue($"Site: {ex.GetType().Name}: {ex.Message}"); } await app.StopAsync(); await minApiTask; Console.WriteLine($"{logs.Count(l => l.StartsWith("Site: Calling Proxy"))} calls from Site -> Api Proxy"); Console.WriteLine($"{logs.Count(l => l.StartsWith("ProxyApiAsync: Calling external"))} calls from Api Proxy -> External"); Console.WriteLine($"{logs.Count(l => l.StartsWith("OperationCanceledException: ProxyApiAsync"))} OperationCanceledException caught in Proxy Api"); Console.WriteLine($"{logs.Count(l => l.StartsWith("OperationCanceledException: Site Delegate"))} OperationCanceledException caught in Site Delegate"); Console.WriteLine($"{logs.Count(l => l.StartsWith("OperationCanceledException: Site ForEachAsync"))} OperationCanceledException caught in Site ForEachAsync Extension"); Console.WriteLine($"{logs.Count(l => l.StartsWith("ProxyApiAsync: ApplicationException"))} ApplicationException caught in Proxy Api"); Console.WriteLine($"{logs.Count(l => l.StartsWith("Site: ApplicationException"))} ApplicationException caught in Site"); Console.WriteLine(""); foreach (var log in logs) { Console.WriteLine(log); } } static async Task WebSiteEndpointAsync( int[] apiInfosToRun, IHttpClientFactory httpClientFactory ) { logs.Enqueue($"Site: {apiInfosToRun.Length} APIs to run"); var apiResponses = await apiInfosToRun.ForEachAsync( new ParallelOptions { MaxDegreeOfParallelism = Int32.MaxValue }, async (item, ct) => { try { logs.Enqueue($"Site: Calling ProxyApiAsync for {item}"); if (item == 3) { // throw new ApplicationException($"Site - Fake Exception for {item}"); } using var client = httpClientFactory.CreateClient(); using var request = new HttpRequestMessage { Method = HttpMethod.Get, RequestUri = new Uri($"http://localhost:5000/proxyapi/{item}"), }; var apiResponse = await client.SendAsync(request, ct); apiResponse.EnsureSuccessStatusCode(); var apiResult = await apiResponse.Content.ReadAsStringAsync(ct); if (apiResult.StartsWith("FAILED")) { throw new ApplicationException(apiResult); } return apiResponse; } catch (OperationCanceledException) { logs.Enqueue($"OperationCanceledException: Site Delegate for {item}"); throw; } } ); logs.Enqueue("Site: Finished all APIs"); } static async Task<string> ProxyApiAsync(int taskId, IHttpClientFactory httpClientFactory, CancellationToken cancellationToken ) { try { logs.Enqueue($"ProxyApiAsync: Calling external API for {taskId}"); var httpClient = httpClientFactory.CreateClient(); var response = await httpClient.GetAsync( "https://www.msn.com/", cancellationToken ); // Simulated call to external api if (taskId == 3) { throw new ApplicationException($"Proxy Api - Fake Exception for {taskId}"); } response.EnsureSuccessStatusCode(); var content = await response.Content.ReadAsStringAsync( cancellationToken ); return $"SUCCESS: Length={content.Length}"; } catch (OperationCanceledException) { logs.Enqueue($"OperationCanceledException: ProxyApiAsync for {taskId}"); return "CANCELLED"; } catch ( Exception ex ) { logs.Enqueue($"ProxyApiAsync: {ex.GetType().Name}: {ex.Message}"); return $"FAILED: {ex.Message}"; } } static async Task<TResult[]> ForEachAsync<TSource, TResult>( this TSource[] source, ParallelOptions parallelOptions, Func<TSource, CancellationToken, ValueTask<TResult>> body ) { TResult[] results = new TResult[source.Length]; await Parallel.ForEachAsync( Enumerable.Range(0, source.Length), parallelOptions, async (i, ct) => { try { results[i] = await body(source[i], ct); // .ConfigureAwait( false ); } catch (OperationCanceledException) { logs.Enqueue($"OperationCanceledException: Site ForEachAsync Extension for {source[i]}"); throw; } }); return results; } }
内容的提问来源于stack exchange,提问作者Terry
相关产品推荐
相关产品推荐

