完整、可直接运行的 C# MQTT 云端通信示例,适用于工业上位机(WinForms / WPF / Worker Service)场景,基于 .NET 8,采用目前最主流、最稳定的 MQTT客户端
·
以下是一个完整、可直接运行的 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();
}));
}
}
}
关键工业级优化点
- 断线自动重连:ManagedMqttClient 内置 AutoReconnectDelay,设置 5s 重试。
- QoS 选择:
- QoS 0:高频非关键数据(如实时温度)
- QoS 1:重要状态(如报警、控制指令)
- QoS 2:关键配置下发(极少用)
- 心跳:订阅 Broker 的 $SYS 主题,或自己每 30s 发一次 ping。
- 断网缓存:用 ConcurrentQueue 缓存未发送数据,重连后批量 Publish。
- 安全:生产环境必须:
- TLS + 证书验证
- 用户名/密码
- 国密 SM4 加密 Payload(可加 BouncyCastle)
- 性能:高频采集(>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)+ 缓冲 + 防丢包完整链路
随时补充!祝项目顺利,数据上云稳稳的!
更多推荐
所有评论(0)