P2: Flexible四类通信锁从全局静态改为按端口/IP粒度(ConcurrentDictionary),不同物理端口互不阻塞

This commit is contained in:
2026-09-14 13:53:46 +08:00
parent 8d8df56544
commit 7043a66a28
4 changed files with 62 additions and 42 deletions
+17 -12
View File
@@ -2,6 +2,7 @@
using NModbus; using NModbus;
using NModbus.Serial; using NModbus.Serial;
using System; using System;
using System.Collections.Concurrent;
using System.IO.Ports; using System.IO.Ports;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
@@ -10,14 +11,18 @@ namespace DeviceCommand.Flexible
{ {
/// <summary> /// <summary>
/// 灵活型 Modbus RTU 串口通信类(免实例化,每次调用临时创建串口连接)。 /// 灵活型 Modbus RTU 串口通信类(免实例化,每次调用临时创建串口连接)。
/// 支持保持寄存器与线圈的读写操作,内部使用通信锁保证同一时刻只有一个事务在执行, /// 支持保持寄存器与线圈的读写操作,内部使用按串口名称粒度的通信锁保证同一串口同一时刻只有一个事务在执行,
/// 适用于偶发性、无需保持长连接的 Modbus RTU 设备读写场景。 /// 不同串口之间互不阻塞,适用于偶发性、无需保持长连接的 Modbus RTU 设备读写场景。
/// </summary> /// </summary>
[ADPCommand] [ADPCommand]
public static class FModbusRTU public static class FModbusRTU
{ {
// 通信锁:保证同一时刻只有一个 Modbus 事务在执行 // 按串口名称粒度的通信锁:同一串口同一时刻只有一个事务在执行,不同串口互不阻塞
private static readonly SemaphoreSlim _commLock = new(1, 1); private static readonly ConcurrentDictionary<string, SemaphoreSlim> _commLocks = new();
/// <summary>获取指定串口的通信锁(不存在则自动创建)</summary>
private static SemaphoreSlim GetLock(string portName)
=> _commLocks.GetOrAdd(portName, _ => new SemaphoreSlim(1, 1));
/// <summary> /// <summary>
/// 创建串口实例并配置超时参数。 /// 创建串口实例并配置超时参数。
@@ -76,7 +81,7 @@ namespace DeviceCommand.Flexible
int writeTimeout = 3000, int writeTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -96,7 +101,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
@@ -127,7 +132,7 @@ namespace DeviceCommand.Flexible
int writeTimeout = 3000, int writeTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -147,7 +152,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
@@ -183,7 +188,7 @@ namespace DeviceCommand.Flexible
int writeTimeout = 3000, int writeTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -203,7 +208,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
@@ -234,7 +239,7 @@ namespace DeviceCommand.Flexible
int writeTimeout = 3000, int writeTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -254,7 +259,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
+17 -12
View File
@@ -1,6 +1,7 @@
using Common.Attributes; using Common.Attributes;
using NModbus; using NModbus;
using System; using System;
using System.Collections.Concurrent;
using System.Net.Sockets; using System.Net.Sockets;
using System.Text; using System.Text;
using System.Threading; using System.Threading;
@@ -10,14 +11,18 @@ namespace DeviceCommand.Flexible
{ {
/// <summary> /// <summary>
/// 灵活型 Modbus TCP 通信类(免实例化,每次调用临时建立 TCP 连接)。 /// 灵活型 Modbus TCP 通信类(免实例化,每次调用临时建立 TCP 连接)。
/// 支持保持寄存器与线圈的读写操作,内部使用通信锁保证同一时刻只有一个事务在执行, /// 支持保持寄存器与线圈的读写操作,内部使用按端点粒度的通信锁保证同一端点同一时刻只有一个事务在执行,
/// 适用于偶发性、无需保持长连接的 Modbus TCP 设备读写场景。 /// 不同端点之间互不阻塞,适用于偶发性、无需保持长连接的 Modbus TCP 设备读写场景。
/// </summary> /// </summary>
[ADPCommand] [ADPCommand]
public static class FModbusTCP public static class FModbusTCP
{ {
// 通信锁:保证同一时刻只有一个 Modbus 事务在执行 // 按端点粒度的通信锁:同一端点同一时刻只有一个事务在执行,不同端点互不阻塞
private static readonly SemaphoreSlim _commLock = new(1, 1); private static readonly ConcurrentDictionary<string, SemaphoreSlim> _commLocks = new();
/// <summary>获取指定端点的通信锁(不存在则自动创建)</summary>
private static SemaphoreSlim GetLock(string ipAddress, int port)
=> _commLocks.GetOrAdd($"{ipAddress}:{port}", _ => new SemaphoreSlim(1, 1));
/// <summary> /// <summary>
/// 建立 TCP 连接并创建 Modbus TCP 主站。 /// 建立 TCP 连接并创建 Modbus TCP 主站。
@@ -57,7 +62,7 @@ namespace DeviceCommand.Flexible
int receiveTimeout = 3000, int receiveTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable; using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable;
@@ -68,7 +73,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
@@ -93,7 +98,7 @@ namespace DeviceCommand.Flexible
int receiveTimeout = 3000, int receiveTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable; using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable;
@@ -103,7 +108,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
@@ -133,7 +138,7 @@ namespace DeviceCommand.Flexible
int receiveTimeout = 3000, int receiveTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable; using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable;
@@ -144,7 +149,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
@@ -169,7 +174,7 @@ namespace DeviceCommand.Flexible
int receiveTimeout = 3000, int receiveTimeout = 3000,
CancellationToken ct = default) CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable; using var master = await ConnectAsync(ipAddress, port, sendTimeout, receiveTimeout, ct) as IDisposable;
@@ -179,7 +184,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
+13 -8
View File
@@ -1,5 +1,6 @@
using Common.Attributes; using Common.Attributes;
using System; using System;
using System.Collections.Concurrent;
using System.IO.Ports; using System.IO.Ports;
using System.Text; using System.Text;
using System.Threading; using System.Threading;
@@ -10,14 +11,18 @@ namespace DeviceCommand.Flexible
/// <summary> /// <summary>
/// 灵活型串口通信类(免实例化,每次调用临时创建串口连接)。 /// 灵活型串口通信类(免实例化,每次调用临时创建串口连接)。
/// 支持只发送指令、发送并读取应答两种最常用操作, /// 支持只发送指令、发送并读取应答两种最常用操作,
/// 内部使用通信锁保证同一时刻只有一个串口事务在执行, /// 内部使用按串口名称粒度的通信锁保证同一串口同一时刻只有一个事务在执行,
/// 适用于偶发性、无需保持长连接的串口设备通信场景(如示波器、电源的 SCPI 指令)。 /// 不同串口之间互不阻塞,适用于偶发性、无需保持长连接的串口设备通信场景(如示波器、电源的 SCPI 指令)。
/// </summary> /// </summary>
[ADPCommand] [ADPCommand]
public static class FSerialPort public static class FSerialPort
{ {
// 通信锁:保证同一时刻只有一个串口事务在执行 // 按串口名称粒度的通信锁:同一串口同一时刻只有一个事务在执行,不同串口互不阻塞
private static readonly SemaphoreSlim _commLock = new(1, 1); private static readonly ConcurrentDictionary<string, SemaphoreSlim> _commLocks = new();
/// <summary>获取指定串口的通信锁(不存在则自动创建)</summary>
private static SemaphoreSlim GetLock(string portName)
=> _commLocks.GetOrAdd(portName, _ => new SemaphoreSlim(1, 1));
/// <summary> /// <summary>
/// 创建串口实例并配置 UTF8 编码与超时参数。 /// 创建串口实例并配置 UTF8 编码与超时参数。
@@ -48,7 +53,7 @@ namespace DeviceCommand.Flexible
/// <param name="ct">异步取消令牌</param> /// <param name="ct">异步取消令牌</param>
public static async Task (string portName,int baudRate,Parity parity,int dataBits,StopBits stopBits,int sendTimeout,int receiveTimeout,string command,CancellationToken ct = default) public static async Task (string portName,int baudRate,Parity parity,int dataBits,StopBits stopBits,int sendTimeout,int receiveTimeout,string command,CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -66,7 +71,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
@@ -90,7 +95,7 @@ namespace DeviceCommand.Flexible
/// <returns>去除结束符并去除首尾空白后的应答字符串</returns> /// <returns>去除结束符并去除首尾空白后的应答字符串</returns>
public static async Task<string> (string portName,int baudRate, Parity parity, int dataBits, StopBits stopBits,int sendTimeout, int receiveTimeout,string command,string delimiter = "\n", CancellationToken ct = default) public static async Task<string> (string portName,int baudRate, Parity parity, int dataBits, StopBits stopBits,int sendTimeout, int receiveTimeout,string command,string delimiter = "\n", CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(portName).WaitAsync(ct);
try try
{ {
using var port = CreatePort( using var port = CreatePort(
@@ -131,7 +136,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(portName).Release();
} }
} }
+15 -10
View File
@@ -1,4 +1,5 @@
using Common.Attributes; using Common.Attributes;
using System.Collections.Concurrent;
using System.Net.Sockets; using System.Net.Sockets;
using System.Text; using System.Text;
@@ -7,14 +8,18 @@ namespace DeviceCommand.Flexible
/// <summary> /// <summary>
/// 灵活型 TCP 通信类(免实例化,每次调用临时建立 TCP 连接)。 /// 灵活型 TCP 通信类(免实例化,每次调用临时建立 TCP 连接)。
/// 支持字节/文本发送、定长字节读取、按结束符读取文本行四种操作, /// 支持字节/文本发送、定长字节读取、按结束符读取文本行四种操作,
/// 内部使用通信锁保证同一时刻只有一个 TCP 事务在执行, /// 内部使用按端点粒度的通信锁保证同一端点同一时刻只有一个 TCP 事务在执行,
/// 适用于偶发性、无需保持长连接的 TCP 设备通信场景。 /// 不同端点之间互不阻塞,适用于偶发性、无需保持长连接的 TCP 设备通信场景。
/// </summary> /// </summary>
[ADPCommand] [ADPCommand]
public static class FTCP public static class FTCP
{ {
// 通信锁:保证同一时刻只有一个 TCP 事务在执行 // 按端点粒度的通信锁:同一端点同一时刻只有一个事务在执行,不同端点互不阻塞
private static readonly SemaphoreSlim _commLock = new(1, 1); private static readonly ConcurrentDictionary<string, SemaphoreSlim> _commLocks = new();
/// <summary>获取指定端点的通信锁(不存在则自动创建)</summary>
private static SemaphoreSlim GetLock(string ipAddress, int port)
=> _commLocks.GetOrAdd($"{ipAddress}:{port}", _ => new SemaphoreSlim(1, 1));
#region Send #region Send
@@ -28,7 +33,7 @@ namespace DeviceCommand.Flexible
/// <param name="ct">异步取消令牌</param> /// <param name="ct">异步取消令牌</param>
public static async Task (string ipAddress,int port,int sendTimeout, byte[] buffer, CancellationToken ct = default) public static async Task (string ipAddress,int port,int sendTimeout, byte[] buffer, CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var client = new TcpClient(); using var client = new TcpClient();
@@ -41,7 +46,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
@@ -73,7 +78,7 @@ namespace DeviceCommand.Flexible
/// <returns>实际读取到的字节数组(对端提前关闭时可能短于请求长度)</returns> /// <returns>实际读取到的字节数组(对端提前关闭时可能短于请求长度)</returns>
public static async Task<byte[]> (string ipAddress,int port,int receiveTimeout,int length,CancellationToken ct = default) public static async Task<byte[]> (string ipAddress,int port,int receiveTimeout,int length,CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var client = new TcpClient(); using var client = new TcpClient();
@@ -100,7 +105,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }
@@ -115,7 +120,7 @@ namespace DeviceCommand.Flexible
/// <returns>去除结束符并去除首尾空白后的文本行</returns> /// <returns>去除结束符并去除首尾空白后的文本行</returns>
public static async Task<string> ( string ipAddress, int port, int receiveTimeout, string delimiter = "\n",CancellationToken ct = default) public static async Task<string> ( string ipAddress, int port, int receiveTimeout, string delimiter = "\n",CancellationToken ct = default)
{ {
await _commLock.WaitAsync(ct); await GetLock(ipAddress, port).WaitAsync(ct);
try try
{ {
using var client = new TcpClient(); using var client = new TcpClient();
@@ -147,7 +152,7 @@ namespace DeviceCommand.Flexible
} }
finally finally
{ {
_commLock.Release(); GetLock(ipAddress, port).Release();
} }
} }