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

.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对象,导致采集目标错误。

适配该场景的实现方案

改用全异步架构,搭配并发度控制,线程利用率会提升数倍,耗时会大幅下降:

  1. 将GetNewReadings方法改造为异步实现,内部HTTP调用使用HttpClient的异步方法,数据库写入也替换为异步API(比如你现有用的Dapper就支持ExecuteAsync等异步方法),方法签名改为async Task GetNewReadingsAsync(RouterModel router, PowerMeterModel meter, CancellationToken stoppingToken)
  2. 使用SemaphoreSlim控制最大并行请求数,避免瞬时请求量过高压垮本机、数据库或者前端设备,可根据服务器配置调整阈值
  3. 收集所有采集任务后用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 10:36:07