探讨 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 通讯设计要素

  1. 连接管理:
    • 使用 ConnectionMultiplexer(StackExchange.Redis)管理连接,自动重连。
    • 配置连接超时和重试策略。
  2. 数据结构:
    • 字符串:存储车牌数据(如 JSON 序列化)。
    • 哈希:存储 PLC 数据字段(状态、置信度、车牌号)。
    • 列表:记录历史车牌数据。
  3. 发布/订阅:
    • 频道命名:如 neuron/plc1/plate,与 MQTT 主题一致。
    • 异步订阅:处理高并发消息。
  4. 重试机制:
    • 集成 Polly 自适应重试(参考前文),处理 Redis 连接中断或超时。
    • 批量重试:合并多个失败操作(如 SET、PUBLISH)。
  5. 幂等性:
    • 写操作(如 SET)使用唯一键或事务 ID。
    • 发布消息记录事务 ID,防止重复处理。
  6. 监控:
    • 记录 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 项目设置

  1. 创建 WinForm 项目:dotnet new winforms -n ModbusRedisIntegration.
  2. 安装 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.
  3. 配置 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.
  4. 配置 Redis:
    • 单机:redis://192.168.10.100:6379,无密码。
    • 集群(可选):配置多个节点,如 192.168.10.100:6379,192.168.10.101:6379.
  5. 配置 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}.
  6. 配置 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 测试步骤

  1. 配置 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.
  2. 配置 Redis:
    • 启动 Redis 服务器(192.168.10.100:6379)。
    • 使用 redis-cli 检查连接:ping(应返回 PONG)。
  3. 配置 EMQX 集群(mTLS,参考前文)。
  4. 运行 QUIC WebSocket 服务端(dotnet run)。
  5. 运行 WinForm 程序,连接 PLC、Modbus UDP、MQTT、QUIC WebSocket、Redis.
  6. 测试功能:
    • 点击“同步数据”,验证多 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(无证书应失败)。
  7. 检查日志,验证 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),请告知!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐