MQTT在印花设备运行状态监控中应用

发布时间:2026/9/4 17:08:39
MQTT在印花设备运行状态监控中应用 一、MQTT原理概述MQTTMessage Queuing Telemetry Transport是一种基于发布/订阅模式的轻量级消息传输协议。其核心通信模型包含三个角色角色说明发布者 (Publisher)发送消息的客户端如印花设备端程序订阅者 (Subscriber)接收消息的客户端如云端监控应用代理 (Broker)位于云端AWS IoT Core的消息中转服务器发布者与订阅者之间不直接通信而是通过Broker进行消息中转。发布者将消息发送到指定的主题TopicBroker负责将消息转发给所有订阅了该主题的客户端。这种模式实现了发布者与订阅者的解耦——发布者无需知道谁在接收消息订阅者也无需知道谁在发布消息。MQTT在印花设备监控中的关键特性轻量级协议头最小仅2字节适合带宽受限的工业现场三种QoS级别可根据不同数据重要性选择可靠性最多一次、至少一次、精确一次遗嘱消息LWT设备异常离线时自动发送告警TLS/SSL加密保证数据传输安全端口8883二、系统总体架构基于MQTT的印花设备运行状态监控系统采用三层架构┌─────────────────────────────────────────────────────────────────┐ │ 第一层数据采集层 │ │ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │ │ │ 印花设备 #1 │ │ 印花设备 #2 │ │ 印花设备 #N │ │ │ │ (C# 发布者) │ │ (C# 发布者) │ │ (C# 发布者) │ │ │ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │ │ │ │ │ │ │ └────────MQTT over TLS (端口8883)────┘ │ └─────────────────────────│──────────────────────────────────────┘ ▼ ┌─────────────────────────────────────────────────────────────────┐ │ 第二层消息代理层 │ │ ┌──────────────────┐ │ │ │ AWS IoT Core │ │ │ │ (MQTT Broker) │ │ │ └────────┬─────────┘ │ │ │ │ │ ┌──────────────┼──────────────┐ │ │ ▼ ▼ ▼ │ │ ┌────────────┐ ┌────────────┐ ┌────────────┐ │ │ │ IoT Rules │ │Device │ │IoT Events │ │ │ │ Engine │ │Shadow │ │(告警) │ │ │ └─────┬──────┘ └─────┬──────┘ └─────┬──────┘ │ │ │ │ │ │ └────────────│──────────────│──────────────│─────────────────────┘ ▼ ▼ ▼ ┌─────────────────────────────────────────────────────────────────┐ │ 第三层数据存储与展示层 │ │ ┌────────────┐ ┌────────────┐ ┌────────────┐ │ │ │ 时序数据库 │ │ 监控大屏 │ │ 移动端App │ │ │ │ (状态存储) │ │ (Web订阅者) │ │ (订阅者) │ │ │ └────────────┘ └────────────┘ └────────────┘ │ └─────────────────────────────────────────────────────────────────┘三、云端实现方案3.1 AWS IoT CoreMQTT BrokerAWS IoT Core作为MQTT Broker负责接收设备消息并转发给订阅者。核心配置步骤创建“事物”Thing为每台印花设备在AWS IoT Core中创建一个Thing代表设备的数字孪生。生成并下载安全凭证设备证书device-certificate.pem.crt私钥device-private.pem.key根CA证书AmazonRootCA1.pem创建IoT策略Policy定义设备允许的操作至少包括{ Effect: Allow, Action: [iot:Connect, iot:Publish, iot:Subscribe, iot:Receive], Resource: [*] }获取终端节点在AWS IoT Core控制台“设置”页面获取终端节点地址如xxxxxxxxxx-ats.iot.region.amazonaws.com3.2 设备影子Device ShadowDevice Shadow是AWS IoT Core的一项服务用于持久化存储设备的最后已知状态即使设备离线也能通过MQTT主题访问其状态。核心MQTT主题以设备IDPrinter01为例主题用途Printer01/shadow/update设备发布状态更新 / 应用发布期望状态Printer01/shadow/update/accepted状态更新被接受的通知Printer01/shadow/update/rejected状态更新被拒绝的通知Printer01/shadow/update/delta期望状态与当前状态存在差异时推送Printer01/shadow/update/documents完整状态文档推送工作原理设备发布状态到shadow/update主题AWS IoT Core将状态写入Device Shadow持久存储应用订阅shadow/update/documents主题获取完整状态变更应用可通过shadow/update发布期望状态如“关机”设备通过订阅shadow/update/delta接收并执行3.3 IoT规则引擎IoT Rules EngineAWS IoT Rules Engine用于对设备消息进行实时处理、转换和路由。典型规则配置示例SELECT deviceId, printerStatus, temperature, timestamp FROM device/status/# WHERE printerStatus Error该规则订阅所有设备状态主题当检测到printerStatus Error时触发后续动作如调用Lambda函数发送告警、写入数据库等。3.4 云端订阅者实现监控应用云端监控应用作为MQTT订阅者订阅设备状态主题以接收实时数据。方案一AWS Lambda 时序数据库Lambda函数订阅设备主题将数据写入Amazon Timestream或InfluxDB监控大屏从数据库读取数据并展示方案二WebSocket连接通过AWS IoT Core的WebSocket支持端口443Web应用可直接订阅MQTT主题实现浏览器端的实时状态监控方案三C# 云端订阅者使用与设备端相同的MQTTnet库在云服务器上运行订阅服务// 订阅设备状态主题 await mqttClient.SubscribeAsync(new MqttTopicFilterBuilder() .WithTopic(device/status/) .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .Build()); // 注册消息接收事件 mqttClient.ApplicationMessageReceived (sender, e) { var payload Encoding.UTF8.GetString(e.ApplicationMessage.Payload); var topic e.ApplicationMessage.Topic; // 解析JSON并存入数据库 / 触发告警 / 更新大屏 Console.WriteLine($收到 {topic}: {payload}); };四、设备端发布者与订阅者实现方案印花设备端C#程序同时扮演两个角色发布者定期发布设备状态、打印机状态等数据订阅者订阅云端下发的控制指令4.1 NuGet包安装Install-Package MQTTnet Install-Package MQTTnet.Extensions.ManagedClient Install-Package Oocx.ReadX509CertificateFromPem4.2 完整代码实现using System; using System.Collections.Generic; using System.IO; using System.Security.Cryptography.X509Certificates; using System.Text; using System.Threading; using System.Threading.Tasks; using MQTTnet; using MQTTnet.Client; using MQTTnet.Extensions.ManagedClient; using Oocx.ReadX509CertificateFromPem; namespace PrinterMqttMonitor { class Program { // 配置区 private const string AwsEndpoint xxxxxxxxxx-ats.iot.region.amazonaws.com; private const int AwsPort 8883; private const string ClientId Printer01; // 与AWS Thing名称一致 // 证书路径建议使用配置文件或环境变量 private const string CaCertPath C:\certs\AmazonRootCA1.pem; private const string DeviceCertPath C:\certs\device-certificate.pem.crt; private const string DeviceKeyPath C:\certs\device-private.pem.key; // 主题定义 private const string TopicStatus device/status; // 发布状态 private const string TopicCommand printer/command; // 订阅指令 private const string TopicShadowUpdate Printer01/shadow/update; // 影子更新 private static IManagedMqttClient _mqttClient; private static readonly Random _random new Random(); static async Task Main(string[] args) { Console.WriteLine( 印花设备MQTT监控客户端启动 ); // 1. 加载证书 var certificates LoadCertificates(); // 2. 创建并配置MQTT客户端 _mqttClient CreateManagedMqttClient(certificates); // 3. 注册消息接收处理作为订阅者 _mqttClient.ApplicationMessageReceived OnMessageReceived; // 4. 启动客户端自动连接与重连 await _mqttClient.StartAsync(CreateClientOptions()); Console.WriteLine(MQTT客户端已启动等待连接...); // 5. 订阅控制指令主题 await _mqttClient.SubscribeAsync(new MqttTopicFilterBuilder() .WithTopic(TopicCommand) .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .Build()); Console.WriteLine($已订阅指令主题: {TopicCommand}); // 6. 主循环定期发布设备状态作为发布者 int sequence 0; while (true) { try { var status CollectDeviceStatus(sequence); await PublishStatusAsync(status); Console.WriteLine($[{DateTime.Now:HH:mm:ss}] 状态已发送); } catch (Exception ex) { Console.WriteLine($发布失败: {ex.Message}); } await Task.Delay(TimeSpan.FromSeconds(30)); // 每30秒发送一次 } } /// summary /// 加载PEM格式证书 /// /summary private static ListX509Certificate LoadCertificates() { // 读取PEM文件 var caPem File.ReadAllText(CaCertPath); var deviceCertPem File.ReadAllText(DeviceCertPath); var deviceKeyPem File.ReadAllText(DeviceKeyPath); // 加载CA证书 var caCert new X509Certificate2(Encoding.UTF8.GetBytes(caPem)); // 使用Oocx库加载带私钥的设备证书 var reader new CertificateFromPemReader(); var deviceCert reader.LoadCertificateWithPrivateKeyFromStrings( deviceCertPem, deviceKeyPem); return new ListX509Certificate { caCert, deviceCert }; } /// summary /// 创建托管MQTT客户端支持自动重连 /// /summary private static IManagedMqttClient CreateManagedMqttClient( ListX509Certificate certificates) { var factory new MqttFactory(); return factory.CreateManagedMqttClient(); } /// summary /// 创建客户端选项 /// /summary private static ManagedMqttClientOptions CreateClientOptions() { var tlsOptions new MqttClientOptionsBuilderTlsParameters { UseTls true, Certificates certificates, // 来自外部变量 SslProtocol System.Security.Authentication.SslProtocols.Tls12, AllowUntrustedCertificates false, IgnoreCertificateChainErrors false, IgnoreCertificateRevocationErrors false }; var options new MqttClientOptionsBuilder() .WithTcpServer(AwsEndpoint, AwsPort) .WithClientId(ClientId) .WithTls(tlsOptions) .WithCleanSession() .WithKeepAlivePeriod(TimeSpan.FromSeconds(60)) // AWS要求30-1200秒 .WithProtocolVersion(MQTTnet.Formatter.MqttProtocolVersion.V311) .Build(); return new ManagedMqttClientOptionsBuilder() .WithClientOptions(options) .WithAutoReconnectDelay(TimeSpan.FromSeconds(5)) .Build(); } /// summary /// 采集设备状态数据 /// /summary private static object CollectDeviceStatus(int sequence) { // 模拟采集印花设备的各种状态 return new { deviceId ClientId, sequence sequence, timestamp DateTime.UtcNow.ToString(o), // 设备状态 deviceStatus Online, uptimeSeconds Environment.TickCount / 1000, // 打印机状态 printerStatus GetPrinterStatus(), inkLevel Math.Round(70 _random.NextDouble() * 25, 1), paperStatus _random.Next(0, 10) 8 ? Low : OK, // 运行参数 temperature Math.Round(35 _random.NextDouble() * 10, 1), currentJobId _random.Next(1000, 9999), jobProgress _random.Next(0, 100), // 故障信息如有 errorCode _random.Next(0, 10) 9 ? E-1023 : null, errorMessage null }; } private static string GetPrinterStatus() { var statuses new[] { Ready, Printing, Paused, Error }; return statuses[_random.Next(statuses.Length)]; } /// summary /// 发布设备状态发布者角色 /// /summary private static async Task PublishStatusAsync(object status) { var json System.Text.Json.JsonSerializer.Serialize(status); var message new MqttApplicationMessageBuilder() .WithTopic(TopicStatus) .WithPayload(Encoding.UTF8.GetBytes(json)) .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) .Build(); await _mqttClient.PublishAsync(message); } /// summary /// 接收云端消息订阅者角色 /// /summary private static void OnMessageReceived(object sender, MqttApplicationMessageReceivedEventArgs e) { var topic e.ApplicationMessage.Topic; var payload Encoding.UTF8.GetString(e.ApplicationMessage.Payload); Console.WriteLine($[指令] 收到主题 {topic}: {payload}); // 解析并执行控制指令 try { var command System.Text.Json.JsonSerializer.DeserializeCommand(payload); ExecuteCommand(command); } catch (Exception ex) { Console.WriteLine($指令解析失败: {ex.Message}); } } private static void ExecuteCommand(Command cmd) { switch (cmd?.Action?.ToLower()) { case start: Console.WriteLine( 执行: 启动打印任务); // 调用打印机启动逻辑 break; case stop: Console.WriteLine( 执行: 停止打印); // 调用打印机停止逻辑 break; case pause: Console.WriteLine( 执行: 暂停打印); break; case resume: Console.WriteLine( 执行: 恢复打印); break; case status_query: Console.WriteLine( 执行: 立即上报状态); // 触发立即发布 break; default: Console.WriteLine($ 未知指令: {cmd?.Action}); break; } } } /// summary /// 控制指令数据结构 /// /summary public class Command { public string Action { get; set; } public string JobId { get; set; } public Dictionarystring, object Parameters { get; set; } } }4.3 关键注意事项AWS不支持遗嘱消息和保留消息如果尝试使用连接会被无提示断开KeepAlive间隔AWS IoT Core要求保持在30-1200秒之间证书加载在C#中直接加载PEM格式证书较为复杂推荐使用Oocx.ReadX509CertificateFromPem库或预先转换为PFX格式使用托管客户端ManagedMqttClient提供自动重连和消息队列管理功能提升稳定性五、实施步骤汇总步骤阶段具体操作1云端准备在AWS IoT Core创建Thing、生成证书、创建策略2云端准备配置IoT规则引擎数据路由、告警触发3云端准备可选部署Device Shadow实现状态持久化4设备端安装NuGet包MQTTnet、Oocx.ReadX509CertificateFromPem5设备端将证书文件部署到设备安全目录6设备端实现证书加载、MQTT连接、状态发布、指令订阅逻辑7设备端集成实际打印机状态采集API替换模拟数据8测试验证使用AWS控制台MQTT测试客户端验证收发9云端部署部署监控大屏应用订阅设备主题展示数据10运维配置CloudWatch告警、设备生命周期事件监控六、主题设计建议主题方向用途device/status/{deviceId}设备→云端定期上报设备与打印机状态device/event/{deviceId}设备→云端突发事件故障、警告即时上报printer/command/{deviceId}云端→设备下发控制指令printer/config/{deviceId}云端→设备下发配置更新$aws/things/{deviceId}/shadow/*双向Device Shadow系统主题