You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Parallel.ForEachAsync与CancellationToken异常不一致问题排查

异步多线程与CancellationToken行为疑问

我对多线程async/await编程经验不足,代码中多处await的耗时不确定,导致日志和行为不一致,因此产生以下疑问。我的系统架构为:网站API端点 → 代理API → 外部API。当前核心困惑是理解CancellationToken的工作机制,以及为何OperationCanceledException的抛出/日志记录不符合预期。我编写了一个.NET 7控制台应用模拟该架构,代码如下。

我的问题

  1. 取消注释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}");行)。原因是什么?
  2. 注释掉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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 16:07:01