.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
相关产品推荐
相关产品推荐

