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

.NET Core BackgroundService中700个MODBUS TCP客户端高效多线程实现

问题描述

基于.NET Core Worker Service模板开发多MODBUS TCP客户端服务,实现了InitializeServices继承自BackgroundService,其ExecuteAsync方法代码如下:

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
    devices = mysqlManager.GetDevices();
    devices.ForEach(x => x.InitializeDevice());
    _exchangeService.SetTags(devices.SelectMany(x => x.ReadTags).ToList());
    devices.ForEach(x => x.Client.InitializeClient());
            
    foreach (var item in devices)
    {        
        Thread thread = new Thread(item.Client.Run);
        thread.Start();
    }

    while (!stoppingToken.IsCancellationRequested)
    {
        Thread.Sleep(10);
        _exchangeService.InsertReadingTagsToDatabase();
    }
}

每个设备对应的Client类包含Run方法,代码如下:

while (true)
{
    try
    {
        using (TcpClient client = new TcpClient())
        {
            client.ReceiveTimeout = this.device.RequestTimeout;
            client.SendTimeout = this.device.RequestTimeout;
            client.SendBufferSize = 100024;
            client.ReceiveBufferSize = 100000;
            client.Client.ReceiveTimeout = this.device.RequestTimeout;
            client.Client.SendTimeout = this.device.RequestTimeout;
            client.Client.SendBufferSize = 100024;
            client.Client.ReceiveBufferSize = 100000;

            if (client.ConnectAsync(ipAdress, portNumber).Wait(this.device.ConnectionTimeout))
            {
                this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 1;
                var factory = new ModbusFactory();
                IModbusMaster master = factory.CreateMaster(client);
                ReadWordsFromDriverServer(master);
                SyncTagsFromRegisters();
                //WriteTagsToDatabase();
                //Console.WriteLine(words.Length);
            }
            else
            {
                this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 0;
                this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").LastReadTime = DateTime.Now;
            }
        }
    }
    catch (Exception ex)
    {}
    Thread.Sleep(100);
}

当前设备数量约700个,需要让700个客户端周期性运行且互不阻塞,但当前代码中客户端互相等待,需定位问题并解决。

问题分析
  • 线程资源过载:创建700个独立Thread,系统线程池默认容量远低于此数,大量线程会导致操作系统上下文切换开销暴增,线程间互相抢占资源,出现等待。
  • 同步阻塞调用:ConnectAsync().Wait()是同步阻塞调用,会占用线程直到连接完成或超时;Thread.Sleep()也会让线程进入阻塞状态,进一步浪费线程资源。
  • 资源重复创建:每次循环都新建TcpClient和ModbusFactory,TCP连接建立/销毁、对象实例化的开销极大,累积后拖慢整体性能。
  • 异常静默吞掉:空catch块隐藏了连接失败、Modbus读写错误等异常,无法排查潜在的阻塞诱因。
  • 数据库插入过于频繁:主循环每10ms调用一次InsertReadingTagsToDatabase(),高频IO操作会占用线程资源,间接影响客户端任务执行。
解决方案

1. 替换线程为异步任务,利用线程池管理

抛弃Thread,改用Task.Run启动异步任务,线程池会自动优化线程数量,避免资源过载:
修改ExecuteAsync方法:

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
    devices = mysqlManager.GetDevices();
    devices.ForEach(x => x.InitializeDevice());
    _exchangeService.SetTags(devices.SelectMany(x => x.ReadTags).ToList());
    devices.ForEach(x => x.Client.InitializeClient());

    // 启动所有设备的异步任务,并绑定取消令牌
    var deviceTasks = devices.Select(item => Task.Run(() => item.Client.RunAsync(stoppingToken), stoppingToken)).ToList();

    // 数据库插入改为异步批量操作,降低频率
    while (!stoppingToken.IsCancellationRequested)
    {
        await Task.Delay(1000, stoppingToken); // 调整为1秒一次,根据实际需求优化
        await _exchangeService.InsertReadingTagsToDatabaseAsync(stoppingToken);
    }

    // 等待所有设备任务完成
    await Task.WhenAll(deviceTasks);
}

2. 将Client.Run改为异步方法,消除阻塞

把Run改成异步方法,用await替代.Wait(),避免线程阻塞:

public async Task RunAsync(CancellationToken stoppingToken)
{
    while (!stoppingToken.IsCancellationRequested)
    {
        try
        {
            using (TcpClient client = new TcpClient())
            {
                client.ReceiveTimeout = this.device.RequestTimeout;
                client.SendTimeout = this.device.RequestTimeout;
                client.SendBufferSize = 100024;
                client.ReceiveBufferSize = 100000;
                client.Client.ReceiveTimeout = this.device.RequestTimeout;
                client.Client.SendTimeout = this.device.RequestTimeout;
                client.Client.SendBufferSize = 100024;
                client.Client.ReceiveBufferSize = 100000;

                // 用异步连接+取消令牌替代Wait
                var connectTask = client.ConnectAsync(ipAdress, portNumber);
                var completedTask = await Task.WhenAny(connectTask, Task.Delay(this.device.ConnectionTimeout, stoppingToken));
                if (completedTask == connectTask)
                {
                    await connectTask; // 确保连接完成,若失败会抛出异常
                    this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 1;
                    var factory = new ModbusFactory();
                    using IModbusMaster master = factory.CreateMaster(client);
                    // 把ReadWordsFromDriverServer改成异步方法
                    await ReadWordsFromDriverServerAsync(master, stoppingToken);
                    SyncTagsFromRegisters();
                }
                else
                {
                    this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 0;
                    this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").LastReadTime = DateTime.Now;
                }
            }
        }
        catch (Exception ex)
        {
            // 记录异常,比如写入日志
            // _logger.LogError(ex, $"设备{this.device.Id}执行失败");
            this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 0;
            this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").LastReadTime = DateTime.Now;
        }
        // 用异步延迟替代Thread.Sleep,不占用线程
        await Task.Delay(100, stoppingToken);
    }
}

3. 优化资源复用(可选)

如果设备支持长连接,可将TcpClient和ModbusMaster移到循环外,仅在连接断开时重新创建,减少连接开销:

public async Task RunAsync(CancellationToken stoppingToken)
{
    TcpClient? client = null;
    IModbusMaster? master = null;
    try
    {
        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                if (client == null || !client.Connected)
                {
                    client?.Dispose();
                    master?.Dispose();
                    
                    client = new TcpClient();
                    // 设置参数...
                    await client.ConnectAsync(ipAdress, portNumber, stoppingToken);
                    var factory = new ModbusFactory();
                    master = factory.CreateMaster(client);
                    this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 1;
                }

                await ReadWordsFromDriverServerAsync(master, stoppingToken);
                SyncTagsFromRegisters();
            }
            catch (Exception ex)
            {
                // 记录异常,标记连接失败
                this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").Value = 0;
                this.device.ReadTags.FirstOrDefault(x => x.Name == "_NoError").LastReadTime = DateTime.Now;
                // 重置连接对象,下次循环重新连接
                client?.Dispose();
                master?.Dispose();
                client = null;
                master = null;
                await Task.Delay(1000, stoppingToken); // 连接失败后延迟1秒重试
            }
            await Task.Delay(100, stoppingToken);
        }
    }
    finally
    {
        client?.Dispose();
        master?.Dispose();
    }
}

4. 优化数据库插入逻辑

将高频的数据库插入改为批量插入,降低IO压力:

  • 修改InsertReadingTagsToDatabase为异步方法InsertReadingTagsToDatabaseAsync
  • 累积一定数量的标签数据后再批量插入,或按固定时间间隔(如1秒)批量写入,避免每10ms一次的频繁IO。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 01:37:16