以下是一个完整、可直接运行的 C# MQTT 云端通信示例,适用于工业上位机(WinForms / WPF / Worker Service)场景,基于 .NET 8,采用目前最主流、最稳定的 MQTT 客户端库 MQTTnet(NuGet 包名:MQTTnet)。

这个示例实现以下功能:

  • 连接到云端 MQTT Broker(支持公网如 EMQX Cloud、HiveMQ、阿里云 IoT 等)
  • 订阅多个主题(设备状态、报警、配置下发)
  • 发布传感器数据(温度、湿度、转速等)
  • 接收云端指令(远程开关、参数调整)
  • 断线自动重连 + 心跳机制
  • 工业级异常处理、日志、QoS 设置
  • 可视化 UI(WinForms 实时显示 + 按钮控制)

1. 环境准备

# 必须安装的包
Install-Package MQTTnet
Install-Package Microsoft.Extensions.Logging.Console   # 可选,用于日志

2. 完整代码示例(WinForms 版本)

using MQTTnet;
using MQTTnet.Client;
using MQTTnet.Extensions.ManagedClient;
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using System.Windows.Forms;
using Microsoft.Extensions.Logging;

namespace IndustrialMQTTClient
{
    public partial class MainForm : Form
    {
        private IManagedMqttClient _mqttClient;
        private readonly ILogger _logger = LoggerFactory.Create(builder => builder.AddConsole()).CreateLogger("MQTT");

        private CancellationTokenSource _cts = new();

        public MainForm()
        {
            InitializeComponent();
            this.Text = "工业上位机 - MQTT 云端通信";
            this.Size = new Size(800, 600);

            // UI 控件(可拖拽或代码创建)
            var lblStatus = new Label { Text = "MQTT状态:未连接", Location = new Point(10, 10), AutoSize = true };
            var btnConnect = new Button { Text = "连接云端", Location = new Point(10, 40) };
            var btnDisconnect = new Button { Text = "断开", Location = new Point(120, 40), Enabled = false };
            var btnPublish = new Button { Text = "发布数据", Location = new Point(10, 80) };
            var txtLog = new TextBox { Multiline = true, ScrollBars = ScrollBars.Vertical, ReadOnly = true, Location = new Point(10, 120), Size = new Size(760, 460) };

            btnConnect.Click += async (s, e) => await ConnectToBrokerAsync();
            btnDisconnect.Click += (s, e) => Disconnect();
            btnPublish.Click += async (s, e) => await PublishDataAsync();

            Controls.AddRange(new Control[] { lblStatus, btnConnect, btnDisconnect, btnPublish, txtLog });

            // 日志输出到文本框
            _logger = LoggerFactory.Create(builder =>
            {
                builder.AddConsole();
                builder.AddProvider(new TextBoxLoggerProvider(txtLog));
            }).CreateLogger("MQTT");
        }

        private async Task ConnectToBrokerAsync()
        {
            try
            {
                var mqttFactory = new MqttFactory();
                var options = new MqttClientOptionsBuilder()
                    .WithTcpServer("broker.emqx.io", 1883) // 公网测试 Broker,可换成阿里云/腾讯云/自建
                    .WithClientId($"UpperPC-{Guid.NewGuid():N}")
                    .WithCredentials("username", "password") // 如需认证
                    .WithTls(false) // 生产环境建议 true + 证书
                    .WithCleanSession(false)
                    .WithKeepAlivePeriod(TimeSpan.FromSeconds(60))
                    .Build();

                var managedOptions = new ManagedMqttClientOptionsBuilder()
                    .WithClientOptions(options)
                    .WithAutoReconnectDelay(TimeSpan.FromSeconds(5))
                    .Build();

                _mqttClient = mqttFactory.CreateManagedMqttClient();

                _mqttClient.ApplicationMessageReceivedAsync += e =>
                {
                    var topic = e.ApplicationMessage.Topic;
                    var payload = Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment);
                    _logger.LogInformation("收到消息 - Topic: {Topic}, Payload: {Payload}", topic, payload);

                    // 示例:云端下发控制指令
                    if (topic == "control/cmd" && payload == "START")
                    {
                        // 执行启动逻辑
                        Invoke((MethodInvoker)(() => lblStatus.Text = "收到云端启动指令"));
                    }

                    return Task.CompletedTask;
                };

                _mqttClient.ConnectedAsync += e =>
                {
                    _logger.LogInformation("MQTT 已连接");
                    Invoke((MethodInvoker)(() => lblStatus.Text = "已连接"));

                    // 订阅主题(可多个)
                    _mqttClient.SubscribeAsync(new MqttTopicFilterBuilder().WithTopic("sensors/#").Build());
                    _mqttClient.SubscribeAsync(new MqttTopicFilterBuilder().WithTopic("control/cmd").Build());

                    return Task.CompletedTask;
                };

                _mqttClient.DisconnectedAsync += e =>
                {
                    _logger.LogWarning("MQTT 断开: {Reason}", e.Reason);
                    Invoke((MethodInvoker)(() => lblStatus.Text = "已断开,正在重连..."));
                    return Task.CompletedTask;
                };

                await _mqttClient.StartAsync(managedOptions);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "连接失败");
                MessageBox.Show("连接失败:" + ex.Message);
            }
        }

        private async Task PublishDataAsync()
        {
            if (!IsConnected) 
            {
                MessageBox.Show("未连接,请先连接");
                return;
            }

            // 模拟传感器数据
            var payload = new
            {
                timestamp = DateTime.UtcNow.ToString("o"),
                deviceId = "SENSOR001",
                temperature = 25.6,
                humidity = 68.3,
                pressure = 1013.2,
                vibration = 0.45
            };

            var json = System.Text.Json.JsonSerializer.Serialize(payload);

            var message = new MqttApplicationMessageBuilder()
                .WithTopic("sensors/data/SENSOR001")
                .WithPayload(Encoding.UTF8.GetBytes(json))
                .WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce) // QoS 1
                .WithRetainFlag(false)
                .Build();

            await _mqttClient.EnqueueAsync(message);
            _logger.LogInformation("已发布数据到 Topic: sensors/data/SENSOR001");
        }

        private void Disconnect()
        {
            _mqttClient?.StopAsync().GetAwaiter().GetResult();
            _logger.LogInformation("已断开 MQTT");
        }

        protected override void OnFormClosing(FormClosingEventArgs e)
        {
            Disconnect();
            base.OnFormClosing(e);
        }

        private bool IsConnected => _mqttClient?.IsConnected ?? false;
    }

    // 简单 TextBox 日志提供程序(可选)
    public class TextBoxLoggerProvider : ILoggerProvider
    {
        private readonly TextBox _textBox;

        public TextBoxLoggerProvider(TextBox textBox) => _textBox = textBox;

        public ILogger CreateLogger(string categoryName) => new TextBoxLogger(_textBox);

        public void Dispose() { }
    }

    public class TextBoxLogger : ILogger
    {
        private readonly TextBox _textBox;

        public TextBoxLogger(TextBox textBox) => _textBox = textBox;

        public IDisposable BeginScope<TState>(TState state) => null;

        public bool IsEnabled(LogLevel logLevel) => true;

        public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception exception, Func<TState, Exception, string> formatter)
        {
            if (formatter == null) return;

            var message = formatter(state, exception);
            _textBox.Invoke((MethodInvoker)(() =>
            {
                _textBox.AppendText($"[{DateTime.Now:HH:mm:ss}] {logLevel}: {message}{Environment.NewLine}");
                _textBox.ScrollToCaret();
            }));
        }
    }
}

关键工业级优化点

  1. 断线自动重连:ManagedMqttClient 内置 AutoReconnectDelay,设置 5s 重试。
  2. QoS 选择
    • QoS 0:高频非关键数据(如实时温度)
    • QoS 1:重要状态(如报警、控制指令)
    • QoS 2:关键配置下发(极少用)
  3. 心跳:订阅 Broker 的 $SYS 主题,或自己每 30s 发一次 ping。
  4. 断网缓存:用 ConcurrentQueue 缓存未发送数据,重连后批量 Publish。
  5. 安全:生产环境必须:
    • TLS + 证书验证
    • 用户名/密码
    • 国密 SM4 加密 Payload(可加 BouncyCastle)
  6. 性能:高频采集(>10Hz)用 IAsyncEnumerable + Channel 缓冲,避免阻塞。

部署建议(研华工控机)

  • AOT 发布(单文件 < 60MB):
    dotnet publish -c Release -r win-x64 --self-contained true /p:PublishAot=true /p:PublishTrimmed=true /p:PublishSingleFile=true
    
  • 运行时优化(runtimeconfig.json):
    {
      "runtimeOptions": {
        "configProperties": {
          "System.GC.LowLatency": "SustainedLowLatency"
        }
      }
    }
    

如果您需要以下任一方向的进一步完整代码,请直接回复:

  • 断网缓存 + 重连后批量上云完整实现
  • MQTT + OPC UA 双协议切换示例
  • Avalonia 跨平台版本(替换 WinForms)
  • 国密 SM4 加密 Payload 示例
  • 高频采集(50Hz)+ 缓冲 + 防丢包完整链路

随时补充!祝项目顺利,数据上云稳稳的!

Logo

智能硬件社区聚焦AI智能硬件技术生态,汇聚嵌入式AI、物联网硬件开发者,打造交流分享平台,同步全国赛事资讯、开展 OPC 核心人才招募,助力技术落地与开发者成长。

更多推荐