.NET Core 3.1环境下多设备并行采集时Parallel.Invoke的替代方案有哪些?
.NET Core 3.1 多设备HTTP采集场景优化方案
现有代码核心问题
你当前使用的Parallel.Invoke是为CPU密集型计算场景设计的API,而你的业务属于典型的IO密集型场景:HTTP请求设备、数据写入数据库都是IO操作,Parallel.Invoke会占用大量线程池线程阻塞等待IO返回,线程池调度开销极高,这是单轮耗时不达标的核心原因。
另外原有代码还存在闭包捕获的潜在bug:循环中直接使用迭代变量创建Action,最终执行时所有Action可能都引用到最后一次迭代的Router和Meter对象,导致采集目标错误。
适配该场景的实现方案
改用全异步架构,搭配并发度控制,线程利用率会提升数倍,耗时会大幅下降:
- 将
GetNewReadings方法改造为异步实现,内部HTTP调用使用HttpClient的异步方法,数据库写入也替换为异步API(比如你现有用的Dapper就支持ExecuteAsync等异步方法),方法签名改为async Task GetNewReadingsAsync(RouterModel router, PowerMeterModel meter, CancellationToken stoppingToken) - 使用
SemaphoreSlim控制最大并行请求数,避免瞬时请求量过高压垮本机、数据库或者前端设备,可根据服务器配置调整阈值 - 收集所有采集任务后用
Task.WhenAll统一等待,全程无阻塞线程,同时间可处理的请求量远高于Parallel.Invoke方案
优化后参考代码
protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // 控制最大并行数,可根据实际情况调整,比如30、50、100 var maxDegreeOfParallelism = 50; var semaphore = new SemaphoreSlim(maxDegreeOfParallelism); while (!stoppingToken.IsCancellationRequested) { _logger.LogInformation("Worker running at: {time}", DateTimeOffset.Now); // Get Relevant Gateways string GetRouters = "SELECT * FROM RouterModel;"; var Routers = await dataAccess.LoadData<RouterModel, dynamic>(GetRouters, new { }, _config.GetConnectionString("GalronDb")); string GetMeters = "SELECT * FROM PowerMeterModel;"; var Meters = await dataAccess.LoadData<PowerMeterModel, dynamic>(GetMeters, new { }, _config.GetConnectionString("GalronDb")); var tasks = new List<Task>(); foreach (RouterModel Router in Routers) { if (Router.IsHavePowerMeters) { foreach (PowerMeterModel Meter in Meters.Where(x => x.IdGateway == Router.Id).ToList()) { if (Meter.IsActive) { // 捕获迭代变量,避免闭包问题 var currentRouter = Router; var currentMeter = Meter; tasks.Add(Task.Run(async () => { await semaphore.WaitAsync(stoppingToken); try { await GetNewReadingsAsync(currentRouter, currentMeter, stoppingToken); } catch (Exception ex) { // 单个设备采集失败不影响全局,这里加日志记录即可 _logger.LogError(ex, "采集设备失败,RouterId:{routerId}, MeterId:{meterId}", currentRouter.Id, currentMeter.Id); } finally { semaphore.Release(); } }, stoppingToken)); } } } } await Task.WhenAll(tasks); _logger.LogInformation("*************** Sync loop is finished ***************"); // 注意原有代码是等待10秒,若需要10分钟请改为 10 * 60 * 1000 await Task.Delay(10 * 60 * 1000, stoppingToken); } }
额外注意点
- HttpClient建议用静态单例或者IHttpClientFactory注入创建,不要每次请求都新建HttpClient,避免套接字资源耗尽
- 如果采集的设备响应速度差异大,可以给每个
GetNewReadingsAsync加超时控制,避免单个慢请求拖慢整体流程 - 数据库写入如果是批量的,可以把采集到的数据先攒到内存队列,批量写入数据库,进一步降低IO开销
内容的提问来源于stack exchange,提问作者Tom Boudniatski
相关产品推荐
相关产品推荐

