探讨 Redis 通讯,并在 C# 中实现 Redis 的 读写 和 订阅/发布功能,结合之前的上下文 Modbus UDP/TCP/RTU、S7.NET、MQTT、mTLS、QUIC WebSock
探讨 Redis 通讯,并在 C# 中实现 Redis 的 读写 和 订阅/发布功能,结合之前的上下文(Modbus UDP/TCP/RTU、S7.NET、MQTT、mTLS、QUIC WebSocket、车牌识别、分布式部署、重试机制、自适应重试、批量重试等),特别是在车牌识别物联网场景中,优化 Modbus UDP 的不可靠性(丢包)和分布式多 PLC 数据采集。
本回答将详细讲解 Redis 通讯的原理、C# 操作 Redis 的实现方法(读写、订阅/发布),并结合前文的自适应和批量重试策略,集成 Redis 用于缓存车牌数据和发布消息。代码基于 .NET 8,使用 StackExchange.Redis 库,支持跨平台(Windows、Linux、macOS),并提供完整的 WinForm 示例。
内容参考搜索结果(如 Redis 协议、StackExchange.Redis)和前文,确保准确性和实用性。
重点优化 Redis 通讯,提供清晰的代码示例。
1. Redis 通讯原理
1.1 Redis 概述Redis(Remote Dictionary Server)是一个高性能、开源的内存键值数据库,支持多种数据结构(如字符串、哈希、列表、集合、有序集合等),广泛用于缓存、消息队列和发布/订阅场景。在车牌识别场景中,Redis 可用于:
- 缓存:存储多 PLC 的车牌数据(寄存器 40108-40110),减少数据库访问。
- 发布/订阅:广播车牌数据或触发拍照指令,替代或补充 MQTT/QUIC WebSocket。
- 分布式支持:通过 Redis 集群或哨兵模式,支持多 PLC 数据同步。
1.2 Redis 通讯原理
- 协议:Redis 使用基于 TCP 的 RESP(REdis Serialization Protocol),轻量高效,支持二进制安全。
- 请求/响应:客户端发送命令(如 SET key value),服务器返回响应(如 +OK)。
- 发布/订阅:客户端订阅频道(SUBSCRIBE channel),发布者向频道发送消息(PUBLISH channel message)。
- 连接:
- 单机:TCP 连接到 Redis 服务器(默认端口 6379)。
- 集群:支持分片和主从复制,客户端通过 MOVED 或 ASK 重定向访问。
- 高性能:
- 内存存储:数据存储在内存,读写速度极快(微秒级)。
- 单线程模型:避免锁竞争,适合高并发场景。
- 持久化:
- RDB(快照):定期保存数据到磁盘。
- AOF(追加文件):记录每次写操作,支持数据恢复。
- 发布/订阅:
- 基于频道的异步消息传递,客户端订阅频道,接收发布者消息。
- 支持模式订阅(如 PSUBSCRIBE pattern*)。
- 事务:
- 使用 MULTI/EXEC 实现原子操作。
- 不支持回滚,需确保幂等性。
1.3 Redis 在车牌识别场景的应用
- 读写:
- 缓存车牌数据:存储 PLC 数据(状态、置信度、车牌号),键格式如 plc:192.168.0.1:plate。
- 控制指令:存储触发拍照状态(线圈 00001),键格式如 plc:192.168.0.1:trigger。
- 分布式同步:多客户端共享 PLC 数据。
- 发布/订阅:
- 广播车牌数据:发布到频道(如 neuron/plc1/plate),客户端(如 Web 界面)订阅。
- 触发拍照:发布指令到频道,PLC 客户端订阅执行。
- 与上下文的联系:
- Modbus UDP:Redis 缓存数据,减少重复采集。
- S7.NET:同步 PLC 数据到 Redis。
- MQTT/QUIC:Redis 发布/订阅可替代或补充 MQTT/QUIC WebSocket。
- 自适应/批量重试:Redis 操作集成重试机制,处理连接中断或超时。
- 分布式部署:Redis 集群支持多 PLC 数据管理。
1.4 Redis 通讯设计要素
- 连接管理:
- 使用 ConnectionMultiplexer(StackExchange.Redis)管理连接,自动重连。
- 配置连接超时和重试策略。
- 数据结构:
- 字符串:存储车牌数据(如 JSON 序列化)。
- 哈希:存储 PLC 数据字段(状态、置信度、车牌号)。
- 列表:记录历史车牌数据。
- 发布/订阅:
- 频道命名:如 neuron/plc1/plate,与 MQTT 主题一致。
- 异步订阅:处理高并发消息。
- 重试机制:
- 集成 Polly 自适应重试(参考前文),处理 Redis 连接中断或超时。
- 批量重试:合并多个失败操作(如 SET、PUBLISH)。
- 幂等性:
- 写操作(如 SET)使用唯一键或事务 ID。
- 发布消息记录事务 ID,防止重复处理。
- 监控:
- 记录 Redis 操作的 RTT、成功率、错误原因。
- 可集成 Prometheus 导出指标。
1.5 Redis 优化思路
- 连接池:ConnectionMultiplexer 单例复用,减少连接开销。
- 批量操作:
- 使用 Batch 或 MULTI 合并多个命令。
- 批量重试失败操作,优化性能。
- 自适应重试:
- 动态调整重试间隔(基于 RTT 滑动平均)。
- 动态调整重试次数(基于成功率)。
- 分布式支持:
- Redis 集群分片存储,负载均衡。
- 哨兵模式确保高可用。
- 跨平台:
- .NET 8 兼容,Redis 客户端支持 Windows、Linux、macOS.
- 确保 Redis 服务器端口(默认 6379)开放。
- 监控集成:
- 使用 INFO 命令监控 Redis 性能。
- 导出 RTT、成功率到 Prometheus。
2. C# 操作 Redis 实现读写和订阅/发布
2.1 实现目标
- 读写功能:
- 存储车牌数据到 Redis(字符串或哈希)。
- 读取 PLC 数据(状态、置信度、车牌号)。
- 实现触发拍照指令存储。
- 订阅/发布功能:
- 发布车牌数据到频道(如 neuron/plc1/plate)。
- 订阅频道,实时接收消息。
- 批量重试:
- 收集 Redis 失败操作,统一重试。
- 按 PLC IP 分组,异步调度。
- 自适应重试:
- 动态调整间隔(RTT 滑动平均,α=0.8)。
- 动态调整次数(成功率 < 50%:5 次,< 80%:4 次,否则 3 次)。
- 分布式支持:
- 管理多 S7-1200(192.168.0.1, 192.168.0.3)。
- Redis 频道隔离(neuron/plc1/plate, neuron/plc2/plate)。
- 功能码:
- Modbus UDP:01(读线圈)、03(读寄存器)、05(写线圈)、06(写寄存器)。
- 集成:
- S7.NET:同步 S7-1200 和 S7-200 SMART。
- MQTT(mTLS):发布车牌数据(与 Redis 互补)。
- QUIC WebSocket(mTLS):广播数据。
- Redis:缓存和发布/订阅。
- 跨平台:.NET 8 兼容 Windows、Linux、macOS.
- 监控:记录 RTT、成功率、批量重试次数。
2.2 项目设置
- 创建 WinForm 项目:dotnet new winforms -n ModbusRedisIntegration.
- 安装 NuGet 包:
- NModbus:dotnet add package NModbus.
- MQTTnet:dotnet add package MQTTnet.
- S7netplus:dotnet add package S7netplus.
- StackExchange.Redis:dotnet add package StackExchange.Redis.
- 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.
- 配置 Redis:
- 单机:redis://192.168.10.100:6379,无密码。
- 集群(可选):配置多个节点,如 192.168.10.100:6379,192.168.10.101:6379.
- 配置 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.Collections.Generic;
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;
using StackExchange.Redis;
namespace ModbusRedisIntegration
{
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 ConcurrentDictionary<string, (double Rtt, int SuccessCount, int TotalCount)> metrics = new();
private readonly Channel<(string Ip, string Operation, Action<IModbusMaster, UdpClient>)> modbusQueue = Channel.CreateUnbounded<(string, string, Action<IModbusMaster, UdpClient>)>();
private readonly Channel<(string, Action<Plc>)> plcQueue = Channel.CreateUnbounded<(string, Action<Plc>)>();
private readonly List<(string Key, string Operation, Action<IDatabase>)> redisFailedRequests = new();
private ConnectionMultiplexer redis;
private IMqttClient mqttClient;
private ClientWebSocket quicWebSocket;
private TextBox textBoxPlcIps, textBoxDb1200, textBoxMqttBroker, textBoxRedis, textBoxLog;
private Button btnConnectPlc, btnConnectModbus, btnConnectMqtt, btnConnectQuic, btnConnectRedis, btnSync, btnTrigger;
private Timer timer;
private ushort transactionId = 0;
public Form1()
{
InitializeComponent();
InitializeUI();
StartRequestProcessors();
StartBatchRetryProcessor();
}
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" };
textBoxRedis = new TextBox { Left = 20, Top = 140, Width = 200, Text = "192.168.10.100:6379" };
btnConnectPlc = new Button { Text = "连接 PLC", Left = 20, Top = 180, Width = 100 };
btnConnectModbus = new Button { Text = "连接 Modbus UDP", Left = 130, Top = 180, Width = 100 };
btnConnectMqtt = new Button { Text = "连接 MQTT", Left = 240, Top = 180, Width = 100 };
btnConnectQuic = new Button { Text = "连接 QUIC", Left = 350, Top = 180, Width = 100 };
btnConnectRedis = new Button { Text = "连接 Redis", Left = 460, Top = 180, Width = 100 };
btnSync = new Button { Text = "同步数据", Left = 570, Top = 180, Width = 100 };
btnTrigger = new Button { Text = "触发拍照", Left = 680, Top = 180, Width = 100 };
textBoxLog = new TextBox { Left = 20, Top = 220, Width = 800, Height = 200, Multiline = true, ScrollBars = ScrollBars.Vertical };
Controls.AddRange(new Control[] { textBoxPlcIps, textBoxDb1200, textBoxMqttBroker, textBoxRedis, btnConnectPlc, btnConnectModbus, btnConnectMqtt, btnConnectQuic, btnConnectRedis, btnSync, btnTrigger, textBoxLog });
btnConnectPlc.Click += BtnConnectPlc_Click;
btnConnectModbus.Click += BtnConnectModbus_Click;
btnConnectMqtt.Click += BtnConnectMqtt_Click;
btnConnectQuic.Click += BtnConnectQuic_Click;
btnConnectRedis.Click += BtnConnectRedis_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);
metrics.TryAdd(ip, (1000, 0, 0));
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("ModbusRedisClient")
.WithWillMessage(new MqttApplicationMessageBuilder()
.WithTopic("/neuron/disconnect")
.WithPayload(Encoding.UTF8.GetBytes("Modbus Redis Client Disconnected"))
.Build())
.Build();
var retryPolicy = CreateAdaptiveRetryPolicy("MQTT");
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 = CreateAdaptiveRetryPolicy("QUIC");
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 void BtnConnectRedis_Click(object sender, EventArgs e)
{
try
{
if (redis == null || !redis.IsConnected)
{
redis = await ConnectionMultiplexer.ConnectAsync(textBoxRedis.Text);
var subscriber = redis.GetSubscriber();
await subscriber.SubscribeAsync("neuron/plc*/plate", (channel, value) =>
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 收到: 频道={channel}, 消息={value}\r\n")));
});
btnConnectRedis.Text = "断开 Redis";
textBoxLog.AppendText($"已连接到 Redis: {textBoxRedis.Text}\r\n");
}
else
{
await redis.CloseAsync();
redis = null;
btnConnectRedis.Text = "连接 Redis";
textBoxLog.AppendText("已断开 Redis\r\n");
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"Redis 连接错误: {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, "WriteCoil", async (master, client) =>
{
var retryPolicy = CreateAdaptiveRetryPolicy(ip);
try
{
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;
UpdateMetrics(ip, responseTime, true);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"触发拍照成功: {ip}, 线圈 00001, 响应时间={responseTime:F0}ms, 成功率={GetSuccessRate(ip):F2}\r\n")));
// Redis 存储触发状态
if (redis?.IsConnected == true)
{
var db = redis.GetDatabase();
var retryPolicyRedis = CreateAdaptiveRetryPolicy("Redis");
try
{
await retryPolicyRedis.ExecuteAsync(async () =>
{
await db.StringSetAsync($"plc:{ip}:trigger", "true", TimeSpan.FromSeconds(60));
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 存储触发状态: plc:{ip}:trigger = true\r\n")));
});
}
catch (Exception ex)
{
lock (redisFailedRequests)
{
redisFailedRequests.Add(($"plc:{ip}:trigger", "StringSet", async db => await db.StringSetAsync($"plc:{ip}:trigger", "true", TimeSpan.FromSeconds(60))));
}
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 存储触发状态失败,加入批量重试: {ip}, 错误={ex.Message}\r\n")));
}
}
});
}
catch (Exception ex)
{
lock (failedRequests)
{
failedRequests.Add((ip, "WriteCoil", async (m, c) => await m.WriteSingleCoilAsync((byte)(ip == "192.168.0.1" ? 17 : 18), 0, true)));
}
UpdateMetrics(ip, 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"触发拍照失败,加入批量重试: {ip}, 错误={ex.Message}\r\n")));
}
}));
}
}
}
catch (Exception ex)
{
textBoxLog.AppendText($"触发拍照错误: {ex.Message}\r\n");
}
}
private async Task SyncDataAsync()
{
try
{
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))
{
await modbusQueue.Writer.WriteAsync((ip, "ReadRegisters", async (master, client) =>
{
var retryPolicy = CreateAdaptiveRetryPolicy(ip);
try
{
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);
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
UpdateMetrics(ip, responseTime, true);
lock (plateDataList)
{
plateDataList.Add((ip, status, confidence, plateNumber));
}
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 解析: {ip}, 状态={status}, 置信度={confidence:F2}, 车牌号={plateNumber}, 响应时间={responseTime:F0}ms, 成功率={GetSuccessRate(ip):F2}\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")));
// Redis 存储车牌数据
if (redis?.IsConnected == true)
{
var db = redis.GetDatabase();
var retryPolicyRedis = CreateAdaptiveRetryPolicy("Redis");
try
{
await retryPolicyRedis.ExecuteAsync(async () =>
{
var plateData = new { Ip = ip, Status = status, Confidence = confidence, PlateNumber = plateNumber, Timestamp = DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss") };
var json = JsonSerializer.Serialize(plateData);
await db.StringSetAsync($"plc:{ip}:plate", json, TimeSpan.FromSeconds(60));
await db.HashSetAsync($"plc:{ip}:data", new[]
{
new HashEntry("status", status),
new HashEntry("confidence", confidence),
new HashEntry("plateNumber", plateNumber)
});
await redis.GetSubscriber().PublishAsync($"neuron/plc{Array.IndexOf(textBoxPlcIps.Text.Split(','), ip) + 1}/plate", json);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 存储: plc:{ip}:plate, 发布到 neuron/plc{Array.IndexOf(textBoxPlcIps.Text.Split(','), ip) + 1}/plate\r\n")));
});
}
catch (Exception ex)
{
lock (redisFailedRequests)
{
var plateData = new { Ip = ip, Status = status, Confidence = confidence, PlateNumber = plateNumber, Timestamp = DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss") };
var json = JsonSerializer.Serialize(plateData);
redisFailedRequests.Add(($"plc:{ip}:plate", "StringSet", async db => await db.StringSetAsync($"plc:{ip}:plate", json, TimeSpan.FromSeconds(60))));
redisFailedRequests.Add(($"plc:{ip}:data", "HashSet", async db => await db.HashSetAsync($"plc:{ip}:data", new[]
{
new HashEntry("status", status),
new HashEntry("confidence", confidence),
new HashEntry("plateNumber", plateNumber)
})));
redisFailedRequests.Add(($"neuron/plc{Array.IndexOf(textBoxPlcIps.Text.Split(','), ip) + 1}/plate", "Publish", async db => await redis.GetSubscriber().PublishAsync($"neuron/plc{Array.IndexOf(textBoxPlcIps.Text.Split(','), ip) + 1}/plate", json)));
}
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 操作失败,加入批量重试: {ip}, 错误={ex.Message}\r\n")));
}
}
await ParseModbusUdpResponseAsync(master, client, unitId, startAddress, 3);
});
}
catch (Exception ex)
{
lock (failedRequests)
{
failedRequests.Add((ip, "ReadRegisters", async (m, c) =>
{
byte unitId = (byte)(ip == "192.168.0.1" ? 17 : 18);
ushort startAddress = 107;
await m.ReadHoldingRegistersAsync(unitId, startAddress, 3);
await m.ReadCoilsAsync(unitId, 0, 1);
await m.WriteSingleRegisterAsync(unitId, startAddress, (ushort)(status + 1));
await ParseModbusUdpResponseAsync(m, c, unitId, startAddress, 3);
}));
}
UpdateMetrics(ip, 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"读取失败,加入批量重试: {ip}, 错误={ex.Message}\r\n")));
}
}));
}
}
await Task.WhenAll(tasks);
// 发布到 MQTT
if (mqttClient?.IsConnected == true)
{
var retryPolicy = CreateAdaptiveRetryPolicy("MQTT");
var failedMessages = new List<(string Topic, MqttApplicationMessage Message)>();
foreach (var data in plateDataList)
{
try
{
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");
});
}
catch (Exception ex)
{
lock (failedMessages)
{
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);
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();
failedMessages.Add((message.Topic, message));
}
textBoxLog.AppendText($"MQTT 发布失败,加入批量重试: {data.Ip}, 错误={ex.Message}\r\n");
}
}
if (failedMessages.Count > 0)
{
await Task.Run(async () =>
{
var retryPolicyMqtt = CreateAdaptiveRetryPolicy("MQTT");
foreach (var (topic, message) in failedMessages)
{
await retryPolicyMqtt.ExecuteAsync(async () =>
{
await mqttClient.PublishAsync(message);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"MQTT 批量重试成功: 主题={topic}\r\n")));
});
}
});
}
}
// 发布到 QUIC WebSocket
if (quicWebSocket?.State == WebSocketState.Open)
{
var retryPolicy = CreateAdaptiveRetryPolicy("QUIC");
var failedMessages = new List<(string Data, ArraySegment<byte> Message)>();
foreach (var data in plateDataList)
{
try
{
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");
});
}
catch (Exception ex)
{
lock (failedMessages)
{
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));
failedMessages.Add((JsonSerializer.Serialize(plateData), new ArraySegment<byte>(message)));
}
textBoxLog.AppendText($"QUIC WebSocket 发布失败,加入批量重试: {data.Ip}, 错误={ex.Message}\r\n");
}
}
if (failedMessages.Count > 0)
{
await Task.Run(async () =>
{
var retryPolicyQuic = CreateAdaptiveRetryPolicy("QUIC");
foreach (var (data, message) in failedMessages)
{
await retryPolicyQuic.ExecuteAsync(async () =>
{
await quicWebSocket.SendAsync(message, WebSocketMessageType.Text, true, CancellationToken.None);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"QUIC WebSocket 批量重试成功: {data}\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 = CreateAdaptiveRetryPolicy("S7.NET");
try
{
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($"S7.NET 同步失败: {ex.Message}\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)
{
var ip = client.Client.RemoteEndPoint.ToString().Split(':')[0];
var retryPolicy = CreateAdaptiveRetryPolicy(ip);
try
{
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;
UpdateMetrics(ip, responseTime, true);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 报文解析: {ip}, TransactionID={receivedTransactionId}, UnitID={receivedUnitId}, 数据={dataStr}, 响应时间={responseTime:F0}ms, 成功率={GetSuccessRate(ip):F2}\r\n")));
});
}
catch (Exception ex)
{
lock (failedRequests)
{
failedRequests.Add((ip, "ParseResponse", async (m, c) => await ParseModbusUdpResponseAsync(m, c, unitId, startAddress, count)));
}
UpdateMetrics(ip, 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus UDP 报文解析失败,加入批量重试: {ip}, 错误={ex.Message}\r\n")));
}
}
private async void StartBatchRetryProcessor()
{
var modbusSemaphore = new SemaphoreSlim(2, 2);
var redisSemaphore = new SemaphoreSlim(2, 2);
_ = Task.Run(async () =>
{
while (true)
{
// Modbus 批量重试
List<(string Ip, string Operation, Action<IModbusMaster, UdpClient>)> modbusBatch;
lock (failedRequests)
{
modbusBatch = failedRequests.OrderBy(r => r.Operation == "WriteCoil" ? 0 : 1).ToList();
failedRequests.Clear();
}
if (modbusBatch.Count > 0)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"开始 Modbus 批量重试: {modbusBatch.Count} 个请求\r\n")));
foreach (var group in modbusBatch.GroupBy(r => r.Ip))
{
await modbusSemaphore.WaitAsync();
try
{
if (modbusMasters.TryGetValue(group.Key, out var master) && udpClients.TryGetValue(group.Key, out var client))
{
var retryPolicy = CreateAdaptiveRetryPolicy(group.Key);
foreach (var request in group)
{
await retryPolicy.ExecuteAsync(async () =>
{
var startTime = DateTime.Now;
await Task.Run(() => request.Item3(master, client));
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
UpdateMetrics(group.Key, responseTime, true);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus 批量重试成功: {group.Key}, 操作={request.Operation}, 响应时间={responseTime:F0}ms\r\n")));
});
}
}
}
catch (Exception ex)
{
UpdateMetrics(group.Key, 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus 批量重试失败: {group.Key}, 错误={ex.Message}\r\n")));
}
finally
{
modbusSemaphore.Release();
}
}
}
// Redis 批量重试
List<(string Key, string Operation, Action<IDatabase>)> redisBatch;
lock (redisFailedRequests)
{
redisBatch = redisFailedRequests.OrderBy(r => r.Operation == "Publish" ? 0 : 1).ToList();
redisFailedRequests.Clear();
}
if (redisBatch.Count > 0 && redis?.IsConnected == true)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"开始 Redis 批量重试: {redisBatch.Count} 个请求\r\n")));
var db = redis.GetDatabase();
foreach (var group in redisBatch.GroupBy(r => r.Key))
{
await redisSemaphore.WaitAsync();
try
{
var retryPolicy = CreateAdaptiveRetryPolicy("Redis");
foreach (var request in group)
{
await retryPolicy.ExecuteAsync(async () =>
{
var startTime = DateTime.Now;
await Task.Run(() => request.Item3(db));
var responseTime = (DateTime.Now - startTime).TotalMilliseconds;
UpdateMetrics("Redis", responseTime, true);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 批量重试成功: {group.Key}, 操作={request.Operation}, 响应时间={responseTime:F0}ms\r\n")));
});
}
}
catch (Exception ex)
{
UpdateMetrics("Redis", 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Redis 批量重试失败: {group.Key}, 错误={ex.Message}\r\n")));
}
finally
{
redisSemaphore.Release();
}
}
}
await Task.Delay(1000);
}
});
}
private IAsyncPolicy CreateAdaptiveRetryPolicy(string key)
{
return Policy
.Handle<Exception>(ex => IsRetriableError(ex))
.WaitAndRetryAsync(
retryCount: CalculateRetryCount(key),
sleepDurationProvider: retryAttempt =>
{
var (rtt, _, _) = metrics.GetOrAdd(key, (1000, 0, 0));
double delay = Math.Min(rtt * Math.Pow(2, retryAttempt), 2000) + new Random().Next(-10, 10);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"重试: {key}, 第{retryAttempt + 1}次, 延时={delay:F0}ms\r\n")));
return TimeSpan.FromMilliseconds(delay);
},
onRetry: (ex, time, retryCount, context) =>
{
UpdateMetrics(key, 0, false);
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"重试: {key}, 第{retryCount}次, 错误={ex.Message}, 成功率={GetSuccessRate(key):F2}\r\n")));
});
}
private bool IsRetriableError(Exception ex)
{
return ex is SocketException || ex is RedisConnectionException || ex.Message.Contains("超时") || ex.Message.Contains("无响应");
}
private int CalculateRetryCount(string key)
{
var (rtt, successCount, totalCount) = metrics.GetOrAdd(key, (1000, 0, 0));
double successRate = totalCount == 0 ? 1.0 : (double)successCount / totalCount;
return successRate < 0.5 ? 5 : successRate < 0.8 ? 4 : 3;
}
private void UpdateMetrics(string key, double responseTime, bool success)
{
metrics.AddOrUpdate(key, (responseTime, success ? 1 : 0, 1), (k, old) =>
{
double newRtt = responseTime > 0 ? old.Rtt * 0.8 + responseTime * 0.2 : old.Rtt;
int newSuccessCount = old.SuccessCount + (success ? 1 : 0);
int newTotalCount = old.TotalCount + 1;
if (newTotalCount > 100)
{
newSuccessCount = (int)(newSuccessCount * 0.5);
newTotalCount = (int)(newTotalCount * 0.5);
}
return (newRtt, newSuccessCount, newTotalCount);
});
}
private double GetSuccessRate(string key)
{
var (rtt, successCount, totalCount) = metrics.GetOrAdd(key, (1000, 0, 0));
return totalCount == 0 ? 1.0 : (double)successCount / totalCount;
}
private async void StartRequestProcessors()
{
var modbusSemaphore = new SemaphoreSlim(2, 2);
_ = 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.Item3(master, client));
}
}
catch (Exception ex)
{
textBoxLog.Invoke((Action)(() => textBoxLog.AppendText(
$"Modbus 请求处理错误: {request.Item1}, 操作={request.Item2}, 错误={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();
redis?.CloseAsync().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 代码解释
- Redis 读写:
- 字符串:StringSetAsync 存储车牌数据 JSON(plc:{ip}:plate),TTL 60 秒。
- 哈希:HashSetAsync 存储字段(plc:{ip}:data,status、confidence、plateNumber)。
- 触发状态:StringSetAsync 存储触发拍照状态(plc:{ip}:trigger)。
- Redis 发布/订阅:
- 发布:PublishAsync 发布车牌数据到 neuron/plc{index}/plate。
- 订阅:SubscribeAsync 订阅 neuron/plc*/plate,实时接收消息。
- 批量重试:
- 失败收集:redisFailedRequests 存储失败操作(StringSet、HashSet、Publish)。
- 批量调度:StartBatchRetryProcessor 每秒检查失败队列,按键分组,优先处理 Publish。
- 并发控制:SemaphoreSlim 限制同时重试 2 个操作。
- 自适应重试:
- RTT 监控:metrics 存储 Redis 和 PLC 的 RTT(滑动平均,α=0.8)、成功次数、总次数。
- 重试间隔:base_delay * 2^retry_attempt,最大 2000ms,±10ms 抖动。
- 重试次数:成功率 < 50%:5 次,< 80%:4 次,否则 3 次。
- 错误过滤:仅重试可恢复错误(RedisConnectionException、超时)。
- 分布式支持:
- 管理多 S7-1200(192.168.0.1, 192.168.0.3),独立 RTT 和成功率。
- Redis 键/频道隔离(plc:{ip}:plate, neuron/plc{index}/plate)。
- Channel 调度 Modbus 和 PLC 请求。
- 功能码:
- 01:读线圈 00001(拍照触发状态)。
- 03:读寄存器 40108-40110(状态、置信度、车牌号)。
- 05:写线圈 00001(触发拍照)。
- 06:写寄存器 40108(状态 +1 测试)。
- S7.NET 集成:
- 同步 S7-1200 和 S7-200 SMART 数据。
- 自适应重试处理 PLC 通信。
- MQTT/QUIC 集成:
- MQTT(mTLS):发布车牌数据,批量重试。
- QUIC WebSocket(mTLS):广播数据,批量重试。
- 跨平台:
- .NET 8 兼容,需开放 UDP 502、TCP 8883、UDP 443、TCP 6379.
- 监控:
- 日志记录 RTT、成功率、批量重试次数、错误原因。
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.
- 配置 Redis:
- 启动 Redis 服务器(192.168.10.100:6379)。
- 使用 redis-cli 检查连接:ping(应返回 PONG)。
- 配置 EMQX 集群(mTLS,参考前文)。
- 运行 QUIC WebSocket 服务端(dotnet run)。
- 运行 WinForm 程序,连接 PLC、Modbus UDP、MQTT、QUIC WebSocket、Redis.
- 测试功能:
- 点击“同步数据”,验证多 PLC 数据读取、Redis/MQTT/QUIC 发布。
- 点击“触发拍照”,验证线圈写入和 Redis 存储(plc:{ip}:trigger)。
- 使用 redis-cli 检查键(GET plc:192.168.0.1:plate)和哈希(HGETALL plc:192.168.0.1:data)。
- 使用 redis-cli monitor 检查发布消息(neuron/plc1/plate)。
- 使用 Wireshark 捕获 UDP 流量(udp.port == 502)。
- 使用 MQTTX 订阅 /neuron/plc1/plate, /neuron/plc2/plate。
- 模拟 Redis 断开(停止服务器),验证批量重试日志(RTT、成功率、次数)。
- 测试 mTLS(无证书应失败)。
- 检查日志,验证 Redis 操作、批量重试、RTT、成功率。
3. Redis 在车牌识别场景的应用
- 场景:
- 多 S7-1200 通过 Modbus UDP 读取车牌数据(寄存器 40108-40110)。
- 触发拍照(线圈 00001)。
- S7-200 SMART 同步数据。
- 数据通过 Redis(缓存/发布)、MQTT(mTLS)、QUIC WebSocket 传输。
- Redis 作用:
- 缓存:减少 PLC 采集频率,降低网络负载。
- 发布/订阅:实时广播车牌数据,替代或补充 MQTT/QUIC。
- 分布式支持:多客户端共享数据,键/频道隔离。
- 可靠性:批量重试确保 Redis 操作成功。
- 优势:
- 高性能:内存存储,微秒级读写。
- 灵活性:支持多种数据结构(字符串、哈希)。
- 实时性:发布/订阅适合车牌数据广播。
- 结合自适应/批量重试,优化网络不稳定场景。
4. 总结
- Redis 通讯原理:
- 基于 RESP 协议,TCP 连接,高性能内存存储。
- 支持读写(字符串、哈希)、发布/订阅、事务。
- C# 实现:
- 读写:StackExchange.Redis 存储车牌数据(StringSetAsync, HashSetAsync),TTL 60 秒。
- 发布/订阅:PublishAsync 发布到 neuron/plc{index}/plate,SubscribeAsync 接收消息。
- 批量重试:redisFailedRequests 收集失败操作,按键分组,优先 Publish。
- 自适应重试:RTT 滑动平均(α=0.8),重试次数 3-5,间隔 100ms-2000ms。
- 分布式:支持多 PLC,Redis 键/频道隔离。
- 集成:S7.NET、MQTT、QUIC WebSocket,mTLS 保护。
- 与上下文的联系:
- Modbus UDP:Redis 缓存减少采集频率,批量重试弥补丢包。
- S7.NET:同步 PLC 数据到 Redis。
- MQTT/QUIC:Redis 发布/订阅互补。
- 车牌识别:可靠缓存和广播数据。
- 跨平台:.NET 8 兼容,需开放 UDP 502、TCP 8883、UDP 443、TCP 6379.
进一步优化建议:
- Redis 集群:部署多节点集群,支持分片和高可用。
- 事务优化:使用 MULTI/EXEC 确保原子性。
- 机器学习:预测 Redis 重试参数。
- Prometheus 集成:导出 Redis RTT、成功率、批量重试指标。
- Lua 脚本:优化复杂 Redis 操作,减少网络往返。
如需实现以上建议或特定场景优化(如 Redis 集群、Lua 脚本、更大规模 PLC),请告知!
更多推荐




所有评论(0)