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

HttpClient.SendAsync高连接限制下仅并行处理两个请求的问题求助

问题:Windows服务并发处理HTTP请求受限(仅2个同时执行)

背景与原始正常方案

我有一个Windows服务,负责从数据库读取数据并通过多个REST API调用处理数据。最初方案通过定时器触发,用SemaphoreSlim限制多线程,数据库读取需等待所有处理完成后才能再次执行,此方案可正常运行:

ServicePointManager.DefaultConnectionLimit = 10;

原始可运行代码:

// 每5秒由定时器触发
private void ProcessTimer_Elapsed(object sender, ElapsedEventArgs e)
{
    var hasLock = false;
    try
    {
        Monitor.TryEnter(timerLock, ref hasLock);
        if (hasLock)
        {
            ProcessNewData();
        }
        else
        {
            log.Info("Failed to acquire lock for timer."); // 经常出现此日志
        }
    }
    finally
    {
        if (hasLock)
        {
            Monitor.Exit(timerLock);
        }
    }
}

public void ProcessNewData()
{
    var unproceesedItems = GetDatabaseItems();

    if (unproceesedItems.Count > 0)
    {
        var downloadTasks = new Task[unproceesedItems.Count];
        var maxThreads = new SemaphoreSlim(semaphoreSlimMinMax, semaphoreSlimMinMax); // semaphoreSlimMinMax = 10 为最大线程数

        for (var i = 0; i < unproceesedItems.Count; i++)
        {
            maxThreads.Wait();
            var iClosure = i;
            downloadTasks[i] =
            Task.Run(async () =>
                {
                    try
                    {
                        await ProcessItemsAsync(unproceesedItems[iClosure]);
                    }
                    catch (Exception ex)
                    {
                        // 异常处理
                    }
                    finally
                    {
                        maxThreads.Release();
                    }
                });
        }

        Task.WaitAll(downloadTasks);
    }
}

重构后的问题

为提升效率,重构后将GetDatabaseItems放到单独线程执行,通过ConcurrentDictionary在取数和处理模块间传递数据(取数模块填充,处理模块清空)。但出现问题:尽管传入10条未处理数据到ProcessItemsAsync,每次仅处理2条,而非设置的10条。

延迟发生在ProcessItemsAsync内部的var response = await client.SendAsync(request);调用处——10个线程均执行到此处,但每次仅2个线程能完成调用。此部分代码在新旧版本中未修改。

新版本修改代码:

public void Start()
{
    ServicePointManager.DefaultConnectionLimit = maxSimultaneousThreads;  // 10

    // 启动未处理数据读取任务
    getUnprocessedDataTimer.Interval = getUnprocessedDataInterval; // 5秒
    getUnprocessedDataTimer.Elapsed += GetUnprocessedData; // 写入ConcurrentDictionary
    getUnprocessedDataTimer.Start();

    cancellationTokenSource = new CancellationTokenSource();

    // 创建新线程处理数据
    Task.Factory.StartNew(() =>
       {
           try
           {
               ProcessNewData(cancellationTokenSource.Token);
           }
           catch (Exception ex)
           {
               // 错误处理
           }
       }, TaskCreationOptions.LongRunning
    );

}

private void ProcessNewData(CancellationToken token)
{
    // 检查任务是否已取消
    while (!token.IsCancellationRequested)
    {
        if (unprocessedDictionary.Count > 0)
        {
            try
            {
                var throttler = new SemaphoreSlim(maxSimultaneousThreads, maxSimultaneousThreads); // maxSimultaneousThreads = 10
                var tasks = unprocessedDictionary.Select(async item =>
                {
                    await throttler.WaitAsync(token);
                    try
                    {
                        if (unprocessedDictionary.TryRemove(item.Key, out var item))
                        {
                            await ProcessItemsAsync(item);
                        }
                    }
                    catch (Exception ex)
                    {
                        // 错误处理
                    }
                    finally
                    {
                        throttler.Release();
                    }
                });
                Task.WhenAll(tasks);
            }
            catch (OperationCanceledException)
            {
                break;
            }
        }

        Thread.Sleep(1000);
    }
}

环境信息

  • .NET Framework 4.7.1
  • Windows Server 2016
  • Visual Studio 2019

已尝试的无效方案

  • 将最大线程数和ServicePointManager.DefaultConnectionLimit设置为30
  • 使用Thread.Start()手动创建线程
  • 替换为同步HttpClient调用
  • 使用Task.Run+Task.WaitAll调用处理逻辑
  • 将长运行线程替换为定时器

需求

实现GetUnprocessedData与ProcessNewData并发执行,确保HttpClient同时处理10个请求。

更新信息

原始项目正常,但添加新定时器或线程后问题复现,移除则恢复正常。


解决方案

1. 修复异步任务等待逻辑

新版本中Task.WhenAll(tasks)未被await,导致循环快速重复执行,多次创建信号量和任务,引发资源竞争。同时需用Task.Delay替代Thread.Sleep,避免阻塞线程池。

修改后的ProcessNewData方法:

private async Task ProcessNewData(CancellationToken token)
{
    while (!token.IsCancellationRequested)
    {
        if (unprocessedDictionary.Count > 0)
        {
            try
            {
                var throttler = new SemaphoreSlim(maxSimultaneousThreads, maxSimultaneousThreads);
                // 批量取出当前所有待处理项,避免遍历过程中字典并发修改
                var itemsToProcess = unprocessedDictionary.ToList();
                // 清空字典,防止后续取数任务重复添加已处理项(根据业务调整,也可逐个移除)
                unprocessedDictionary.Clear();

                var tasks = itemsToProcess.Select(async item =>
                {
                    await throttler.WaitAsync(token);
                    try
                    {
                        await ProcessItemsAsync(item.Value);
                    }
                    catch (Exception ex)
                    {
                        // 异常处理
                    }
                    finally
                    {
                        throttler.Release();
                    }
                });
                await Task.WhenAll(tasks); // 必须等待所有任务完成
            }
            catch (OperationCanceledException)
            {
                break;
            }
        }

        await Task.Delay(1000, token); // 异步等待,不阻塞线程
    }
}

启动代码对应调整为异步:

// 创建新线程处理数据
_ = Task.Factory.StartNew(async () =>
   {
       try
       {
           await ProcessNewData(cancellationTokenSource.Token);
       }
       catch (Exception ex)
       {
           // 错误处理
       }
   }, cancellationTokenSource.Token, TaskCreationOptions.LongRunning, TaskScheduler.Default);

2. 确保HttpClient复用与连接池配置

问题核心原因是每次创建新的HttpClient实例,导致默认每个域名的连接数限制(2个)生效,覆盖了ServicePointManager.DefaultConnectionLimit的设置。必须复用单例HttpClient:

// 类级别静态单例HttpClient,避免频繁创建实例
private static readonly HttpClient _httpClient = new HttpClient();

async Task ProcessItemsAsync(Item item)
{
    var request = new HttpRequestMessage(HttpMethod.Post, "你的API地址");
    // 设置请求内容等
    var response = await _httpClient.SendAsync(request);
    // 处理响应
}

同时确保ServicePointManager.DefaultConnectionLimit在所有HttpClient实例创建前设置,建议放在服务启动的最开头:

public void Start()
{
    // 必须在创建HttpClient前设置
    ServicePointManager.DefaultConnectionLimit = maxSimultaneousThreads;
    // 其他启动逻辑...
}

3. 优化ConcurrentDictionary访问逻辑

避免在遍历字典的同时执行TryRemove,改为批量取出所有待处理项后再处理,防止并发修改导致的遍历异常或重复操作。

4. 定时器线程优化

确保GetUnprocessedData方法仅负责快速读取数据库并写入ConcurrentDictionary,不要在定时器的Elapsed事件中执行长时间操作,避免阻塞定时器线程池。


内容的提问来源于stack exchange,提问作者eglease

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 18:40:58