结合之前的上下文(Modbus UDP/TCP/RTU、S7.NET、MQTT、mTLS、QUIC WebSocket、车牌识别、分布式部署、功能码扩展等),特别是针对 Modbus UDP 的不可靠
结合之前的上下文(Modbus UDP/TCP/RTU、S7.NET、MQTT、mTLS、QUIC WebSocket、车牌识别、分布式部署、功能码扩展等),特别是针对 Modbus UDP 的不可靠性(UDP 丢包)。本回答将详细讲解重试机制的原理、设计策略、优化方法,以及在 C# WinForm 中实现 Modbus UDP 重试机制的完整示例。代码基于 .NET 8,集成 S7.NET 和 MQTT(mTLS),针对车牌识别物联网场景,跨平台兼容(Windows、Linux、macOS),并参考搜索结果(如 Polly 库、Modbus 协议)和前文内容,确保准确性和实用性。
1. 重试机制原理
1.1 重试机制概述重试机制是一种在通信或操作失败时自动重新尝试的策略,广泛用于处理网络不稳定、超时或临时故障,尤其在 Modbus UDP 场景中,因 UDP 不保证数据到达,必须在应用层实现可靠性。重试机制的目标是提高通信成功率,同时避免过度重试导致资源浪费或网络拥堵。
- 适用场景:
- Modbus UDP:UDP 无连接,数据包可能丢失或乱序。
- MQTT:网络中断可能导致发布失败。
- QUIC/WebSocket:连接断开或消息丢失。
- S7.NET:PLC 通信超时或连接中断。
- 车牌识别:确保车牌数据可靠传输到云端。
- 关键问题:
- 如何检测失败(超时、异常码、丢包)。
- 确定重试次数和间隔。
- 避免无限重试或网络拥堵。
- 确保重试逻辑高效且易于维护。
1.2 重试机制核心原理
- 失败检测:
- 超时:设置接收超时(如 1000ms),无响应视为失败。
- 异常码:Modbus 返回异常码(如 0x02 非法地址,0x03 非法值)。
- 网络错误:Socket 异常(如 SocketException)。
- 自定义验证:检查响应数据完整性(如事务标识符不匹配)。
- 重试策略:
- 固定重试:固定间隔(如 100ms)重试固定次数(如 3 次)。
- 指数退避:重试间隔随次数递增(如 100ms、200ms、400ms),减少拥堵。
- 自适应重试:根据网络状态动态调整间隔和次数。
- 随机抖动:在间隔中添加随机延迟(如 ±10ms),避免重试风暴。
- 终止条件:
- 达到最大重试次数。
- 操作成功(如收到有效响应)。
- 遇到不可恢复的错误(如设备离线)。
- 事务管理:
- Modbus UDP 使用事务标识符(Transaction ID)区分请求/响应对。
- 每次重试递增事务标识符,确保响应匹配。
1.3 重试机制设计要素
- 超时配置:
- 设置合理的超时(如 1000ms),根据设备响应时间调整。
- UDP 通常需要较短超时,因其延迟低。
- 重试次数:
- 典型值为 3-5 次,平衡可靠性和性能。
- 过多重试可能导致延迟或资源耗尽。
- 重试间隔:
- 固定间隔:简单,但可能引发拥堵。
- 指数退避:间隔递增,适合网络不稳定场景。
- 随机抖动:避免多个客户端同时重试。
- 错误分类:
- 可重试错误:超时、丢包、临时网络中断。
- 不可重试错误:非法功能码、设备离线、证书错误。
- 日志记录:
- 记录每次重试的错误原因、时间戳和上下文。
- 便于调试和性能分析。
- 幂等性:
- 确保重试操作不产生副作用(如重复写入寄存器)。
- Modbus UDP 读操作天然幂等,写操作需验证响应。
1.4 与上下文的联系
- Modbus UDP(前文):
- UDP 无内置重传,重试机制弥补可靠性不足。
- 重试处理超时和丢包,结合事务标识符验证响应。
- Modbus TCP/RTU(前文):
- TCP 由协议栈保证可靠性,重试需求较低。
- RTU 使用 CRC 校验,重试处理串口错误。
- S7.NET(前文):
- 重试机制可用于 PLC 通信超时。
- MQTT(mTLS,前文):
- 重试确保 MQTT 发布成功,结合 QoS 1/2。
- QUIC/WebSocket(前文):
- 重试处理 WebSocket 消息丢失或连接断开。
- 分布式部署(前文):
- 重试机制支持多 PLC 并发通信,动态调整策略。
- 车牌识别:
- 确保车牌数据(寄存器 40108-40110)可靠传输。
- 触发拍照(线圈 00001)需幂等重试。
1.5 优化重试机制
- Polly 库:
- 使用 Polly 实现重试、指数退避和抖动。
- 支持异步操作,集成 .NET 8。
- 自适应重试:
- 根据网络延迟动态调整超时和间隔。
- 使用滑动窗口统计成功率,优化重试次数。
- 批量重试:
- 针对分布式 PLC,批量重试失败请求,减少开销。
- 优先级调度:
- 关键操作(如触发拍照)优先重试。
- 监控:
- 记录重试次数、延迟和失败率,集成 Prometheus 或日志系统。
- 幂等性保障:
- 写操作(如功能码 05/06)检查响应,防止重复执行。
1.6 跨平台注意事项
- .NET 8:支持 Polly、UdpClient,兼容 Windows、Linux、macOS。
- 超时配置:Linux 可能需要调整 Socket 超时(SO_RCVTIMEO)。
- 日志:使用 Serilog 或 Microsoft.Extensions.Logging 跨平台记录。
- 端口:开放 UDP 502(Modbus)、TCP 8883(MQTT)、UDP 443(QUIC)。
2. C# WinForm 实现优化重试机制
2.1 实现目标
- 重试机制:
- 使用 Polly 实现指数退避和随机抖动。
- 支持 Modbus UDP 读/写操作(功能码 01、03、05、06)。
- 自适应重试:根据响应时间动态调整间隔。
- 记录重试日志(时间、错误、PLC IP)。
- 分布式支持:
- 管理多个 S7-1200(Modbus UDP 从站)。
- MQTT 主题隔离(/neuron/plc1/plate、/neuron/plc2/plate)。
- 功能码:
- 01(读线圈)、03(读寄存器)、05(写线圈)、06(写寄存器)。
- 集成:
- S7.NET:同步 S7-1200 和 S7-200 SMART。
- MQTT(mTLS):发布车牌数据。
- QUIC WebSocket(mTLS):广播数据。
- 跨平台:.NET 8 兼容 Windows、Linux、macOS。
2.2 项目设置
- 创建 WinForm 项目:dotnet new winforms -n ModbusUdpRetry.
- 安装 NuGet 包:
- NModbus:dotnet add package NModbus.
- MQTTnet:dotnet add package MQTTnet.
- S7netplus:dotnet add package S7netplus.
- Polly:dotnet add package Polly.
- 配置 PLC:
- S7-1200(PLC1):IP 192.168.0.1,Unit ID 17,寄存器 40108-40110(状态、置信度、车牌号),线圈 00001(触发拍照)。
- S7-1200(PLC2):IP 192.168.0.3,Unit ID 18,寄存器 40108-40110。
- S7-200 SMART:IP 192.168.0.2,DB1。
- 配置 EMQX 集群(参考前文 mTLS):
- 地址:mqtts://192.168.10.174:8883.
- 证书:ca.crt, client.pfx(密码:password)。
- ACL:conf
{allow, {user, "ModbusClient"}, publish, ["/neuron/plc1/plate", "/neuron/plc2/plate"]}. {allow, {user, "ModbusClient"}, subscribe, ["/neuron/#"]}. {deny, all}.
- 配置 QUIC WebSocket 服务端(参考前文,wss://localhost:5001/ws/plate,mTLS)。
2.3 代码示例
2.3.1 WinForm 客户端csharp
using System;
using System.Collections.Concurrent;
using System.Net;
using System.Net.Security;
using System.Net.Sockets;
using System.Net.WebSockets;
using System.Security.Cryptography.X509Certificates;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;
using System.Windows.Forms;
using MQTTnet;
using MQTTnet.Client;
using NModbus;
using Polly;
using S7.Net;
namespace ModbusUdpRetry
{
public partial class Form1 : Form
{
private readonly ConcurrentDictionary<string, Plc> plcs = new();
private readonly ConcurrentDictionary<string, IModbusMaster> modbusMasters = new();
private readonly ConcurrentDictionary<string, UdpClient> udpClients = new();
private readonly Channel<(string, Action<IModbusMaster, UdpClient>)> modbusQueue = Channel.CreateUnbounded<(string, Action<IModbusMaster, UdpClient>)>();
private readonly Channel<(string, Action<Plc>)> plcQueue = Channel.CreateUnbounded<(string, Action<Plc>)>();
private IMqttClient mqttClient;
private ClientWebSocket quicWebSocket;
private TextBox textBoxPlcIps, textBoxDb1200, textBoxMqttBroker, textBoxLog;
private Button btnConnectPlc, btnConnectModbus, btnConnectMqtt, btnConnectQuic, btnSync, btnTrigger;
private Timer timer;
private ushort transactionId = 0;
private readonly ConcurrentDictionary<string, double> responseTimes = new(); // 记录 PLC 响应时间
public Form1()
{
InitializeComponent();
InitializeUI();
StartRequestProcessors();
}
private void InitializeUI()
{
textBoxPlcIps = new TextBox { Left = 20, Top = 20, Width = 200, Text = "192.168.0.1,192.168.0.3" };
textBoxDb1200 = new TextBox { Left = 20, Top = 60, Width = 200, Text = "1" };
textBoxMqttBroker = new TextBox { Left = 20, Top = 100, Width = 200, Text = "mqtts://192.168.10.174:8883" };
btnConnectPlc = new Button { Text = "连接 PLC", Left = 20, Top = 140, Width = 100 };
btnConnectModbus = new Button { Text = "连接 Modbus UDP", Left = 130, Top = 140, Width = 100 };
btnConnectMqtt = new Button { Text = "连接 MQTT", Left = 240, Top = 140, Width = 100 };
btnConnectQuic = new Button { Text = "连接 QUIC", Left = 350, Top = 140, Width = 100 };
btnSync = new Button { Text = "同步数据", Left = 460, Top = 140, Width = 100 };
btnTrigger = new Button { Text = "触发拍照", Left = 570, Top = 140, Width = 100 };
textBoxLog = new TextBox { Left = 20, Top = 180, Width = 700, Height = 200, Multiline = true, ScrollBars = ScrollBars.Vertical };
Controls.AddRange(new Control[] { textBoxPlcIps, textBoxDb1200, textBoxMqttBroker, btnConnectPlc, btnConnectModbus, btnConnectMqtt, btnConnectQuic, btnSync, btnTrigger, textBoxLog });
btnConnectPlc.Click += BtnConnectPlc_Click;
btnConnectModbus.Click += BtnConnectModbus_Click;
btnConnectMqtt.Click += BtnConnectMqtt_Click;
btnConnectQuic.Click += BtnConnectQuic_Click;
btnSync.Click += BtnSync_Click;
btnTrigger.Click += BtnTrigger_Click;
timer = new Timer { Interval = 2000 };
timer.Tick += async (s, e) => await SyncDataAsync();
}
private async void BtnConnectPlc_Click(object sender, EventArgs e)
{
try
{
if (!plcs.ContainsKey("S7-1200"))
{
var plc1200 = new Plc(CpuType.S71200, textBoxPlcIps.Text.Split(',')[0], 0, 0);
plcs.TryAdd("S7-1200", plc1200);
await Task.Run(() => plc1200.Open());
if (plc1200.IsConnected)
textBoxLog.AppendText($"已连接到 S7-1200: {textBoxPlcIps.Text.Split(',')[0]}\r\n");
else
textBoxLog.AppendText("S7-1200 连接失败\r\n");
}
if (!plcs.ContainsKey("S7-200 SMART"))
{
var plc200Smart = new Plc(CpuType.S71200, "192.168.0.2", 0, 0);
plcs.TryAdd("S7-200 SMART", plc200Smart);
await Task.Run(() => plc200Smart.Open());
if (plc200Smart.IsConnected)
textBoxLog.AppendText($"已连接到 S7-200 SMART: 192.168.0.2\r\n");
else
textBoxLog.AppendText("S7-200 SMART 连接失败\r\n");
}
btnConnectPlc.Text = plcs.Count > 0 ? "断开 PLC" : "连接 PLC";
timer.Enabled = plcs.Count > 0;
}
catch (Exception ex)
{
textBoxLog.AppendText($"PLC 连接错误: {ex.Message}\r\n");
}
}
private async void BtnConnectModbus_Click(object sender, EventArgs e)
{
try
{
var ips = textBoxPlcIps.Text.Split(',');
foreach (var ip in ips)
{
if (!modbusMasters.ContainsKey(ip))
{
var factory = new ModbusFactory();
var udpClient = new UdpClient();
udpClient.Connect(ip, 502);
udpClient.Client.ReceiveTimeout = 1000;
var modbusMaster = factory.CreateMaster(udpClient);
modbusMasters.TryAdd(ip, modbusMaster);
udpClients.TryAdd(ip, udpClient);
responseTimes.TryAdd(ip, 1000); // 初始响应时间 1000ms
textBoxLog.AppendText($"已连接到 Modbus UDP: {ip}:502\r\n");
}
}
btnConnectModbus.Text = modbusMasters.Count > 0 ? "断开 Modbus UDP" : "连接 Modbus UDP";
}
catch (Exception ex)
{
textBoxLog.AppendText($"Modbus UDP 连接错误: {ex.Message}\r\n");
}
}
private async void BtnConnectMqtt_Click(object sender, EventArgs e)
{
try
{
if (mqttClient == null || !mqttClient.IsConnected)
{
var factory = new MqttFactory();
mqttClient = factory.CreateMqttClient();
var clientCertificate = new X509Certificate2("client.pfx", "password");
var caCertificate = new X509Certificate2("ca.crt");
var options = new MqttClientOptionsBuilder()
.WithTcpServer("192.168.10.174", 8883)
.WithTls(tlsOptions =>
{
tlsOptions.SslProtocol = System.Security.Authentication.SslProtocols.Tls13;
tlsOptions.Certificates = new[] { clientCertificate };
tlsOptions.ClientCertificatesCallback = certs => new[] { clientCertificate };
tlsOptions.CertificateValidationHandler = context =>
{
var chain = new X509Chain();
chain.ChainPolicy.ExtraStore.Add(caCertificate);
chain.ChainPolicy.VerificationFlags = X509VerificationFlags.AllowUnknownCertificateAuthority;
return chain.Build(context.Certificate);
};
})
.WithClientId("ModbusUdpClient")
.WithWillMessage(new MqttApplicationMessageBuilder()
.WithTopic("/neuron/disconnect")
.WithPayload(Encoding.UTF8.GetBytes("Modbus UDP Client Disconnected"))
.Build())
.Build();
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100 + new Random().Next(-10, 10)),
(ex, time) => textBoxLog.AppendText($"MQTT 连接重试: 错误={ex.Message}\r\n"));
await retryPolicy.ExecuteAsync(async () => await mqttClient.ConnectAsync(options));
btnConnectMqtt.Text = "断开 MQTT";
textBoxLog.AppendText($"已连接到 MQTT Broker (mTLS): {textBoxMqttBroker.Text}\r\n");
await mqttClient.SubscribeAsync(new MqttTopicFilterBuilder().WithTopic("/neuron/#").Build());
mqttClient.ApplicationMessageReceivedAsync += e =>
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"MQTT 收到: 主题={e.ApplicationMessage.Topic}, 消息={Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment)}\r\n")));
return Task.CompletedTask;
};
}
else
{
await mqttClient.DisconnectAsync();
mqttClient = null;
btnConnectMqtt.Text = "连接 MQTT";
textBoxLog.AppendText("已断开 MQTT\r\n");
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"MQTT 连接错误: {ex.Message}\r\n");
}
}
private async void BtnConnectQuic_Click(object sender, EventArgs e)
{
try
{
if (quicWebSocket == null || quicWebSocket.State != WebSocketState.Open)
{
quicWebSocket = new ClientWebSocket();
quicWebSocket.Options.AddSubProtocol("websocket");
quicWebSocket.Options.ClientCertificates = new X509CertificateCollection { new X509Certificate2("client.pfx", "password") };
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100 + new Random().Next(-10, 10)),
(ex, time) => textBoxLog.AppendText($"QUIC 连接重试: 错误={ex.Message}\r\n"));
await retryPolicy.ExecuteAsync(async () => await quicWebSocket.ConnectAsync(new Uri("wss://localhost:5001/ws/plate"), CancellationToken.None));
btnConnectQuic.Text = "断开 QUIC";
textBoxLog.AppendText("已连接到 QUIC WebSocket: wss://localhost:5001/ws/plate\r\n");
_ = Task.Run(() => ReceiveQuicMessagesAsync(quicWebSocket));
}
else
{
await quicWebSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "客户端关闭", CancellationToken.None);
quicWebSocket = null;
btnConnectQuic.Text = "连接 QUIC";
textBoxLog.AppendText("已断开 QUIC WebSocket\r\n");
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"QUIC WebSocket 错误: {ex.Message}\r\n");
}
}
private async Task ReceiveQuicMessagesAsync(ClientWebSocket webSocket)
{
var buffer = new byte[1024 * 4];
try
{
while (webSocket.State == WebSocketState.Open)
{
var result = await webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None);
if (result.MessageType == WebSocketMessageType.Text)
{
string message = Encoding.UTF8.GetString(buffer, 0, result.Count);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"QUIC WebSocket 收到: {message}\r\n")));
}
else if (result.MessageType == WebSocketMessageType.Close)
{
await webSocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "服务端关闭", CancellationToken.None);
break;
}
}
}
catch (Exception ex)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"QUIC WebSocket 接收错误: {ex.Message}\r\n")));
}
}
private async void BtnSync_Click(object sender, EventArgs e)
{
await SyncDataAsync();
}
private async void BtnTrigger_Click(object sender, EventArgs e)
{
try
{
foreach (var ip in textBoxPlcIps.Text.Split(','))
{
if (modbusMasters.TryGetValue(ip, out var modbusMaster) && udpClients.TryGetValue(ip, out var udpClient))
{
await modbusQueue.Writer.WriteAsync((ip, async (master, client) =>
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt =>
{
double baseDelay = responseTimes.GetOrAdd(ip, 1000);
double delay = Math.Min(baseDelay * Math.Pow(2, retryAttempt), 2000) + new Random().Next(-10, 10);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"触发拍照重试: {ip}, 延时={delay:F0}ms\r\n")));
return TimeSpan.FromMilliseconds(delay);
}, (ex, time, retryCount, context) =>
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"触发拍照重试: {ip}, 第{retryCount}次, 错误={ex.Message}\r\n")));
});
await retryPolicy.ExecuteAsync(async () =>
{
var startTime = DateTime.Now;
await master.WriteSingleCoilAsync((byte)(ip == "192.168.0.1" ? 17 : 18), 0, true);
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
responseTimes[ip] = responseTimes.GetOrAdd(ip, 1000) * 0.8 + responseTime * 0.2; // 自适应调整
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"触发拍照成功: {ip}, 线圈 00001, 响应时间={responseTime:F0}ms\r\n")));
});
}));
}
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"触发拍照错误: {ex.Message}\r\n");
}
}
private async Task SyncDataAsync()
{
try
{
// Modbus UDP: 读取多个 PLC 数据
var tasks = new List<Task>();
var plateDataList = new List<(string Ip, int Status, float Confidence, string PlateNumber)>();
foreach (var ip in textBoxPlcIps.Text.Split(','))
{
if (modbusMasters.TryGetValue(ip, out var modbusMaster) && udpClients.TryGetValue(ip, out var udpClient))
{
tasks.Add(modbusQueue.Writer.WriteAsync((ip, async (master, client) =>
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt =>
{
double baseDelay = responseTimes.GetOrAdd(ip, 1000);
double delay = Math.Min(baseDelay * Math.Pow(2, retryAttempt), 2000) + new Random().Next(-10, 10);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"读取重试: {ip}, 延时={delay:F0}ms\r\n")));
return TimeSpan.FromMilliseconds(delay);
}, (ex, time, retryCount, context) =>
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"读取重试: {ip}, 第{retryCount}次, 错误={ex.Message}\r\n")));
});
await retryPolicy.ExecuteAsync(async () =>
{
var startTime = DateTime.Now;
byte unitId = (byte)(ip == "192.168.0.1" ? 17 : 18);
ushort startAddress = 107; // 40108 - 40001 = 107
// 读保持寄存器 (功能码 03)
ushort[] registers = await master.ReadHoldingRegistersAsync(unitId, startAddress, 3);
byte[] rawData = new byte[6];
for (int i = 0; i < 3; i++)
{
rawData[i * 2] = (byte)(registers[i] >> 8);
rawData[i * 2 + 1] = (byte)(registers[i] & 0xFF);
}
var status = registers[0];
var confidence = BitConverter.ToSingle(new byte[] { rawData[3], rawData[2], rawData[5], rawData[4] }, 0);
var plateNumber = Encoding.ASCII.GetString(rawData, 4, 2);
lock (plateDataList)
{
plateDataList.Add((ip, status, confidence, plateNumber));
}
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
responseTimes[ip] = responseTimes.GetOrAdd(ip, 1000) * 0.8 + responseTime * 0.2; // 自适应调整
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 解析: {ip}, 状态={status}, 置信度={confidence:F2}, 车牌号={plateNumber}, 响应时间={responseTime:F0}ms\r\n")));
// 读线圈 (功能码 01)
bool[] coils = await master.ReadCoilsAsync(unitId, 0, 1);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 线圈: {ip}, 线圈 00001={coils[0]}\r\n")));
// 写寄存器 (功能码 06)
await master.WriteSingleRegisterAsync(unitId, startAddress, (ushort)(status + 1));
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 写入寄存器: {ip}, 地址={startAddress}, 值={status + 1}\r\n")));
await ParseModbusUdpResponseAsync(master, client, unitId, startAddress, 3);
});
})));
}
}
await Task.WhenAll(tasks);
// 发布到 MQTT
if (mqttClient?.IsConnected == true)
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100 + new Random().Next(-10, 10)),
(ex, time) => textBoxLog.AppendText($"MQTT 发布重试: 错误={ex.Message}\r\n"));
foreach (var data in plateDataList)
{
await retryPolicy.ExecuteAsync(async () =>
{
var plateData = new { data.Ip, data.PlateNumber, data.Confidence, data.Status, Timestamp = DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss") };
var json = JsonSerializer.Serialize(plateData);
if (json.Length > 1024)
{
textBoxLog.AppendText($"MQTT 消息过大,取消发布: {data.Ip}\r\n");
return;
}
var message = new MqttApplicationMessageBuilder()
.WithTopic($"/neuron/plc{Array.IndexOf(textBoxPlcIps.Text.Split(','), data.Ip) + 1}/plate")
.WithPayload(Encoding.UTF8.GetBytes(json))
.WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce)
.WithRetainFlag(true)
.Build();
await mqttClient.PublishAsync(message);
textBoxLog.AppendText($"MQTT 发布: {json}\r\n");
});
}
}
// 发布到 QUIC WebSocket
if (quicWebSocket?.State == WebSocketState.Open)
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100 + new Random().Next(-10, 10)),
(ex, time) => textBoxLog.AppendText($"QUIC 发布重试: 错误={ex.Message}\r\n"));
foreach (var data in plateDataList)
{
await retryPolicy.ExecuteAsync(async () =>
{
var plateData = new { data.Ip, data.PlateNumber, data.Confidence, data.Status, Timestamp = DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss") };
var message = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(plateData));
await quicWebSocket.SendAsync(new ArraySegment<byte>(message), WebSocketMessageType.Text, true, CancellationToken.None);
textBoxLog.AppendText($"QUIC WebSocket 发布: {JsonSerializer.Serialize(plateData)}\r\n");
});
}
}
// S7.NET 同步
if (plcs.TryGetValue("S7-1200", out var plc1200) && plcs.TryGetValue("S7-200 SMART", out var plc200Smart) && plc1200.IsConnected && plc200Smart.IsConnected)
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromMilliseconds(Math.Pow(2, retryAttempt) * 100 + new Random().Next(-10, 10)),
(ex, time) => textBoxLog.AppendText($"S7.NET 重试: 错误={ex.Message}\r\n"));
await retryPolicy.ExecuteAsync(async () =>
{
int db1200 = int.Parse(textBoxDb1200.Text);
var readTasks = new[]
{
Task.Run(() => plc1200.Read(DataType.DataBlock, db1200, 0, VarType.Byte, 1)),
Task.Run(() => plc1200.Read(DataType.DataBlock, db1200, 4, VarType.Real, 1)),
Task.Run(() => plc1200.Read(DataType.DataBlock, db1200, 10, VarType.String, 10))
};
var results = await Task.WhenAll(readTasks);
var status = Convert.ToByte(results[0]);
var confidence = Convert.ToSingle(results[1]);
var plateNumber = results[2].ToString();
await Task.Run(() =>
{
plc200Smart.Write(DataType.DataBlock, db1200, 0, status);
plc200Smart.Write(DataType.DataBlock, db1200, 4, confidence);
plc200Smart.Write(DataType.DataBlock, db1200, 10, plateNumber);
});
textBoxLog.AppendText($"S7.NET 同步: 状态={status}, 置信度={confidence:F2}, 车牌号={plateNumber}\r\n");
});
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"同步错误: {ex.Message}\r\n");
}
}
private async Task ParseModbusUdpResponseAsync(IModbusMaster master, UdpClient client, byte unitId, ushort startAddress, int count)
{
try
{
var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryAsync(3, retryAttempt =>
{
double baseDelay = responseTimes.GetOrAdd(client.Client.RemoteEndPoint.ToString().Split(':')[0], 1000);
double delay = Math.Min(baseDelay * Math.Pow(2, retryAttempt), 2000) + new Random().Next(-10, 10);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"报文解析重试: 延时={delay:F0}ms\r\n")));
return TimeSpan.FromMilliseconds(delay);
}, (ex, time, retryCount, context) =>
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"报文解析重试: 第{retryCount}次, 错误={ex.Message}\r\n")));
});
await retryPolicy.ExecuteAsync(async () =>
{
var startTime = DateTime.Now;
byte[] request = new byte[] {
(byte)(transactionId >> 8), (byte)(transactionId & 0xFF),
0x00, 0x00,
0x00, 0x06,
unitId,
0x03,
(byte)(startAddress >> 8), (byte)(startAddress & 0xFF),
(byte)(count >> 8), (byte)(count & 0xFF)
};
transactionId++;
await client.SendAsync(request, request.Length);
var result = await client.ReceiveAsync();
byte[] response = result.Buffer;
if (response.Length < 9) throw new Exception("响应长度不足");
ushort receivedTransactionId = (ushort)((response[0] << 8) | response[1]);
ushort protocolId = (ushort)((response[2] << 8) | response[3]);
ushort length = (ushort)((response[4] << 8) | response[5]);
byte receivedUnitId = response[6];
byte functionCode = response[7];
byte byteCount = response[8];
if (protocolId != 0 || receivedUnitId != unitId || functionCode != 0x03)
throw new Exception("无效响应");
string dataStr = "";
for (int i = 0; i < byteCount / 2; i++)
{
ushort value = (ushort)((response[9 + i * 2] << 8) | response[9 + i * 2 + 1]);
dataStr += $"寄存器 {startAddress + i}: {value}, ";
}
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
responseTimes[client.Client.RemoteEndPoint.ToString().Split(':')[0]] = responseTimes.GetOrAdd(client.Client.RemoteEndPoint.ToString().Split(':')[0], 1000) * 0.8 + responseTime * 0.2;
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 报文解析: TransactionID={receivedTransactionId}, UnitID={receivedUnitId}, 数据={dataStr}, 响应时间={responseTime:F0}ms\r\n")));
});
}
catch (Exception ex)
{
textBoxLog.AppendText($"Modbus UDP 报文解析错误: {ex.Message}\r\n");
}
}
private async void StartRequestProcessors()
{
var modbusSemaphore = new SemaphoreSlim(1, 1);
_ = Task.Run(async () =>
{
while (await modbusQueue.Reader.WaitToReadAsync())
{
if (modbusQueue.Reader.TryRead(out var request))
{
await modbusSemaphore.WaitAsync();
try
{
if (modbusMasters.TryGetValue(request.Item1, out var master) && udpClients.TryGetValue(request.Item1, out var client))
{
await Task.Run(() => request.Item2(master, client));
}
}
catch (Exception ex)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"Modbus 请求处理错误: {request.Item1}, {ex.Message}\r\n")));
}
finally
{
modbusSemaphore.Release();
}
}
}
});
var plcSemaphore = new SemaphoreSlim(1, 1);
while (await plcQueue.Reader.WaitToReadAsync())
{
if (plcQueue.Reader.TryRead(out var request))
{
await plcSemaphore.WaitAsync();
try
{
if (plcs.TryGetValue(request.Item1, out var plc) && plc.IsConnected)
{
await Task.Run(() => request.Item2(plc));
}
}
catch (Exception ex)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText($"PLC 请求处理错误: {ex.Message}\r\n")));
}
finally
{
plcSemaphore.Release();
}
}
}
}
protected override void OnFormClosing(FormClosingEventArgs e)
{
timer?.Stop();
foreach (var plc in plcs.Values)
plc.Close();
foreach (var client in udpClients.Values)
client.Close();
mqttClient?.DisconnectAsync().GetAwaiter().GetResult();
quicWebSocket?.CloseAsync(WebSocketCloseStatus.NormalClosure, "窗体关闭", CancellationToken.None).GetAwaiter().GetResult();
base.OnFormClosing(e);
}
}
static class Program
{
[STAThread]
static void Main()
{
Application.EnableVisualStyles();
Application.SetCompatibleTextRenderingDefault(false);
Application.Run(new Form1());
}
}
}
2.3.2 QUIC WebSocket 服务端复用前文 mTLS 配置的 QUIC WebSocket 服务端代码(wss://localhost:5001/ws/plate),支持 mTLS,广播 Modbus UDP 数据。
2.4 代码解释
- 重试机制:
- Polly 实现:3 次重试,指数退避(100ms、200ms、400ms,最大 2000ms),添加 ±10ms 随机抖动。
- 自适应调整:记录每个 PLC 的响应时间(responseTimes),使用加权平均(80% 历史 + 20% 当前)动态调整基线延迟。
- 应用范围:
- Modbus UDP:读寄存器(03)、读线圈(01)、写线圈(05)、写寄存器(06)、报文解析。
- MQTT:连接和发布。
- QUIC WebSocket:连接和发送。
- S7.NET:PLC 数据读写。
- 日志:记录重试次数、错误原因、延迟和响应时间。
- 分布式支持:
- 管理多个 S7-1200(192.168.0.1, 192.168.0.3),使用 ConcurrentDictionary 存储 UdpClient 和 IModbusMaster。
- MQTT 主题隔离(/neuron/plc1/plate, /neuron/plc2/plate)。
- 使用 Channel 和 SemaphoreSlim 调度并发请求。
- 功能码:
- 01:读线圈 00001(拍照触发状态)。
- 03:读寄存器 40108-40110(状态、置信度、车牌号)。
- 05:写线圈 00001(触发拍照,按钮触发)。
- 06:写寄存器 40108(状态 +1 测试)。
- 报文解析:
- 构造 MBAP + PDU,发送 UDP 数据报。
- 解析响应,验证事务标识符、Unit ID、功能码。
- 重试机制确保解析可靠性。
- S7.NET 集成:
- 同步 S7-1200 和 S7-200 SMART 数据。
- 使用 Polly 重试 PLC 通信。
- MQTT 集成:
- mTLS 连接 EMQX,发布分布式 PLC 数据。
- 重试确保发布成功,QoS 1。
- QUIC WebSocket:
- 广播数据,mTLS 保护,重试发送。
- 跨平台:
- .NET 8 兼容,UDP 502、TCP 8883、UDP 443 端口需开放。
- 错误处理:
- 捕获超时、Socket 异常、Modbus 异常码。
- 区分可重试(超时、丢包)和不可重试错误(非法功能码)。
2.5 测试步骤
- 配置 PLC:
- S7-1200(PLC1):IP 192.168.0.1,Unit ID 17,寄存器 40108-40110,线圈 00001。
- S7-1200(PLC2):IP 192.168.0.3,Unit ID 18,寄存器 40108-40110。
- S7-200 SMART:IP 192.168.0.2,DB1。
- 配置 EMQX 集群(mTLS,参考前文)。
- 运行 QUIC WebSocket 服务端(dotnet run)。
- 运行 WinForm 程序,连接 PLC、Modbus UDP、MQTT 和 QUIC WebSocket。
- 测试功能:
- 点击“同步数据”,验证多 PLC 数据读取、MQTT/QUIC 发布。
- 点击“触发拍照”,验证线圈写入(功能码 05)。
- 使用 Wireshark 捕获 UDP 流量(udp.port == 502)。
- 使用 MQTTX 订阅 /neuron/plc1/plate 和 /neuron/plc2/plate。
- 模拟丢包(关闭 PLC),验证重试日志(延时、次数、错误)。
- 测试 mTLS(无证书应失败)。
- 检查日志,验证响应时间、成功率和重试次数。
3. 重试机制在车牌识别场景的应用
- 场景:
- S7-1200 通过 Modbus UDP 读取车牌数据(寄存器 40108-40110)。
- 触发拍照(线圈 00001)。
- S7-200 SMART 同步数据。
- 数据通过 MQTT(mTLS)或 QUIC WebSocket 传输。
- 重试机制作用:
- 可靠性:确保车牌数据在网络不稳定时成功传输。
- 实时性:自适应调整延迟,优化响应时间。
- 幂等性:写线圈/寄存器验证响应,防止重复操作。
- 分布式支持:多 PLC 数据采集,动态重试。
- 优势:
- 弥补 UDP 不可靠性,接近 TCP 可靠性。
- 自适应重试减少网络拥堵。
- 详细日志便于调试。
4. 总结
- 重试机制原理:
- 检测超时、丢包、异常码,执行指数退避和随机抖动。
- 自适应调整:根据响应时间动态优化延迟。
- 使用 Polly 实现异步重试,记录日志。
- C# 实现:
- Polly:3 次重试,指数退避(100ms-2000ms,±10ms 抖动)。
- 自适应:滑动平均调整基线延迟。
- 应用:Modbus UDP(01、03、05、06)、MQTT、QUIC、S7.NET。
- 分布式:支持多 PLC,主题隔离,异步调度。
- 与上下文的联系:
- Modbus UDP:重试弥补丢包,支持多功能码。
- S7.NET:重试处理 PLC 通信超时。
- MQTT:重试确保发布成功,mTLS 保护。
- QUIC/WebSocket:重试发送广播消息。
- 车牌识别:可靠采集和传输车牌数据。
- 跨平台:.NET 8 兼容,需开放 UDP 502、TCP 8883、UDP 443。
进一步优化建议:
- 高级自适应:使用机器学习预测最佳重试间隔。
- 监控集成:导出重试数据到 Prometheus,分析成功率。
- 批量重试:优化多 PLC 请求,减少重试开销。
- 幂等性增强:为写操作添加事务 ID,防止重复执行。
如需实现以上建议或特定场景优化(如大规模 PLC 集群、自定义重试策略),请告知!
更多推荐


所有评论(0)