HttpClient.SendAsync高连接限制下仅并行处理两个请求的问题求助
背景与原始正常方案
我有一个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

