-
C#网络编程之MQTT编程(MQTTnet、EMQX)
第61章 C# MQTT编程(MQTTnet、EMQX)
一、我踩过的MQTT坑:从“M2Mqtt崩溃到凌晨3点”到“MQTTnet+EMQX救了我”
做智能家居平台的时候,一开始用老的M2Mqtt库做客户端,结果上线后遇到高vb.net教程C#教程python教程SQL教程access 2010教程并发(1000个设备同时连接)直接崩溃了——查了半天发现M2Mqtt不支持.NET Core异步编程,线程池被耗尽。后来换成MQTTnet库,异步编程+连接池,瞬间支撑了10万级并发。还有一次用公共MQTT Broker测试,结果消息被别人订阅了,导致用户隐私泄露——后来自己部署EMQX集群,配置ACL权限,再也没出过安全问题。这节我把这些踩坑经验揉进去,用大白话讲透MQTTnet的C#实战和EMQX的生产级部署,代码逐行拆解,拓展生产级优化技巧,让你在C# MQTT开发中少走弯路!
二、MQTTnet:C# MQTT开发的“瑞士军刀”,异步高性能首选
MQTTnet是.NET生态最流行的MQTT库,核心是“全异步编程、跨平台、功能丰富、高性能”,支持.NET Framework、.NET Core、.NET 5+,适合从简单设备到大型平台的所有MQTT场景。
大白话解释:把MQTTnet比作“C# MQTT的万能工具箱”
1.异步编程:用async/await处理连接、发布、订阅,不阻塞主线程,性能拉满;
2.跨平台:支持Windows、Linux、macOS,甚至嵌入式设备(比如树莓派);
3.功能丰富:支持MQTT 3.1.1/5.0、QoS 0/1/2、遗嘱消息、保留消息、会话持久化;
4.高性能:单线程每秒处理10万+消息,支持连接池、批量消息。
我踩过的坑:一开始用MQTTnet的同步方法,结果设备连接时主线程被阻塞,UI界面卡死——后来换成异步方法,UI丝滑流畅!
实战1:C# MQTTnet客户端(发布+订阅+遗嘱消息)
步骤1:安装NuGet包
bash
# 核心库
Install-Package MQTTnet
# 客户端库(可选,核心库已包含)
Install-Package MQTTnet.Extensions.ManagedClient
步骤2:MQTTnet客户端完整代码(逐行讲解)
csharp
using System;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using MQTTnet;
using MQTTnet.Client;
using MQTTnet.Extensions.ManagedClient;
namespace CSharpMQTT.MQTTnetDemo;
public class MqttNetClient
{
// 托管客户端:自动处理重连、消息队列、会话持久化
private IManagedMqttClient _managedClient;
private readonly string _brokerAddress = "localhost"; // 自己部署的EMQX地址
private readonly int _brokerPort = 1883;
private readonly string _clientId = $"home_client_{Guid.NewGuid():N}";
private readonly string _username = "home_user"; // EMQX创建的用户名
private readonly string _password = "Home@123"; // EMQX创建的密码
// 初始化客户端
public async Task InitializeAsync()
{
try
{
// 1. 创建MQTT工厂
var factory = new MqttFactory();
// 2. 创建托管客户端实例
_managedClient = factory.CreateManagedMqttClient();
// 3. 配置托管客户端选项
var managedOptions = new ManagedMqttClientOptionsBuilder()
.WithAutoReconnectDelay(TimeSpan.FromSeconds(5)) // 自动重连间隔5秒
.WithClientOptions(CreateMqttClientOptions()) // 基础客户端选项
.Build();
// 4. 注册事件回调
_managedClient.ConnectedAsync += OnConnectedAsync;
_managedClient.DisconnectedAsync += OnDisconnectedAsync;
_managedClient.ApplicationMessageReceivedAsync += OnMessageReceivedAsync;
// 5. 启动托管客户端
await _managedClient.StartAsync(managedOptions);
Console.WriteLine("MQTTnet托管客户端初始化完成");
}
catch (Exception ex)
{
Console.WriteLine($"客户端初始化失败:{ex.Message}");
}
}
// 创建基础MQTT客户端选项
private MqttClientOptions CreateMqttClientOptions()
{
// 遗嘱消息:客户端异常断开时,Broker自动发布此消息
var willMessage = new MqttApplicationMessageBuilder()
.WithTopic("home/clients/status")
.WithPayload(Encoding.UTF8.GetBytes($"{_clientId} disconnected"))
.WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce)
.WithRetainFlag(true) // 保留消息,新订阅者能收到
.Build();
// 客户端连接选项
return new MqttClientOptionsBuilder()
.WithTcpServer(_brokerAddress, _brokerPort)
.WithClientId(_clientId)
.WithCredentials(_username, _password) // 用户名密码认证
.WithWillMessage(willMessage) // 设置遗嘱消息
.WithCleanSession(false) // 持久化会话,重连后接收离线消息
.WithKeepAlivePeriod(TimeSpan.FromSeconds(60)) // 心跳间隔60秒
.Build();
}
// 发布消息
public async Task PublishMessageAsync(string topic, string payload, bool retain = false)
{
if (!_managedClient.IsConnected)
{
Console.WriteLine("客户端未连接,无法发布消息");
return;
}
try
{
var message = new MqttApplicationMessageBuilder()
.WithTopic(topic)
.WithPayload(Encoding.UTF8.GetBytes(payload))
.WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce)
.WithRetainFlag(retain)
.Build();
// 托管客户端会自动处理消息队列,离线时缓存消息,上线后发送
await _managedClient.EnqueueAsync(message);
Console.WriteLine($"发布消息成功:主题={topic},内容={payload}");
}
catch (Exception ex)
{
Console.WriteLine($"发布消息失败:{ex.Message}");
}
}
// 订阅主题
public async Task SubscribeTopicAsync(string topic)
{
if (!_managedClient.IsConnected)
{
Console.WriteLine("客户端未连接,无法订阅主题");
return;
}
try
{
var subscription = new MqttTopicFilterBuilder()
.WithTopic(topic)
.WithQualityOfServiceLevel(MQTTnet.Protocol.MqttQualityOfServiceLevel.AtLeastOnce)
.Build();
await _managedClient.SubscribeAsync(subscription);
Console.WriteLine($"订阅主题成功:{topic}");
}
catch (Exception ex)
{
Console.WriteLine($"订阅主题失败:{ex.Message}");
}
}
// 断开连接
public async Task DisconnectAsync()
{
if (_managedClient.IsConnected)
{
await _managedClient.StopAsync();
Console.WriteLine("客户端已断开连接");
}
}
#region 事件回调
// 连接成功回调
private Task OnConnectedAsync(MqttClientConnectedEventArgs args)
{
Console.WriteLine($"客户端连接成功:ClientId={_clientId},Broker={_brokerAddress}:{_brokerPort}");
return Task.CompletedTask;
}
// 断开连接回调
private Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs args)
{
string reason = args.Exception != null ? args.Exception.Message : "正常断开";
Console.WriteLine($"客户端断开连接:原因={reason},将在5秒后自动重连");
return Task.CompletedTask;
}
// 收到消息回调
private Task OnMessageReceivedAsync(MqttApplicationMessageReceivedEventArgs args)
{
string topic = args.ApplicationMessage.Topic;
string payload = Encoding.UTF8.GetString(args.ApplicationMessage.Payload);
var qos = args.ApplicationMessage.QualityOfServiceLevel;
Console.WriteLine($"收到消息:主题={topic},内容={payload},QoS={qos}");
// 模拟处理控制指令:比如收到"on"就打开灯光
if (topic == "home/livingroom/light/control" && payload == "on")
{
Console.WriteLine("执行操作:打开客厅灯光");
}
return Task.CompletedTask;
}
#endregion
}
// 测试代码:模拟智能家居场景
class Program
{
static async Task Main(string[] args)
{
var mqttClient = new MqttNetClient();
await mqttClient.InitializeAsync();
// 订阅控制指令主题
await mqttClient.SubscribeTopicAsync("home/livingroom/light/control");
// 模拟每2秒上报一次温度
var cts = new CancellationTokenSource();
_ = Task.Run(async () =>
{
int i = 0;
while (!cts.Token.IsCancellationRequested)
{
string temperature = $"25.{i}°C";
await mqttClient.PublishMessageAsync("home/livingroom/temperature", temperature, retain: true);
i = (i + 1) % 10;
await Task.Delay(2000);
}
}, cts.Token);
Console.WriteLine("按任意键退出...");
Console.ReadKey();
cts.Cancel();
await mqttClient.DisconnectAsync();
}
}
逐行拆解核心代码:
1.托管客户端:IManagedMqttClient是MQTTnet的核心,自动处理重连、离线消息缓存、会话持久化,不用自己写复杂的重连逻辑;
2.遗嘱消息:客户端异常断开时,Broker自动发布状态消息,比如通知其他设备“客厅温度传感器离线了”;
3.持久化会话:WithCleanSession(false)表示重连后Broker会发送离线期间的消息,适合需要接收所有消息的场景;
4.消息队列:EnqueueAsync会把消息加入队列,离线时缓存,上线后自动发送,避免消息丢失;
5.事件回调:用异步事件处理连接、断开、消息接收,不阻塞主线程,适合UI或高并发场景。
三、EMQX:生产级MQTT Broker的“扛把子”,高可用集群部署
EMQX是开源的高性能MQTT Broker,核心是“百万级并发、分布式集群、全链路监控、安全可靠”,适合生产环境的物联网平台、智能家居、车联网等场景。
大白话解释:把EMQX比作“物联网消息快递站”
1.Broker集群:多个快递站协同工作,处理百万级包裹(消息);
2.ACL权限:只有授权用户才能存取自己的快递(消息);
3.监控面板:实时查看快递站的吞吐量、延迟、异常;
4.规则引擎:自动分拣快递,比如把温度消息转发到数据库,把控制指令转发到设备。
我踩过的坑:一开始用单节点EMQX,结果Broker挂了导致所有设备离线——后来部署EMQX集群,用负载均衡,即使一个节点挂了,其他节点继续工作,高可用拉满!
实战2:EMQX生产级部署与使用
步骤1:Docker快速部署EMQX(适合测试和生产)
bash
# 拉取EMQX 5.x镜像
docker pull emqx/emqx:5.6.1
# 启动单节点EMQX,映射端口
docker run -d --name emqx -p 1883:1883 -p 8083:8083 -p 8084:8084 -p 18083:18083 emqx/emqx:5.6.1
1883:MQTT TCP端口;
8083:MQTT WebSocket端口;
8084:MQTT WebSocket TLS端口;
18083:EMQX控制台端口。
步骤2:EMQX控制台使用(创建用户、配置ACL)
1.打开浏览器访问http://localhost:18083,默认用户名admin,密码public;
2.创建用户:进入“访问控制”→“用户”→“创建用户”,比如用户名home_user,密码Home@123;
3.配置ACL规则:进入“访问控制”→“ACL”→“创建规则”,比如允许home_user订阅和发布home/#主题;
4.查看监控:进入“监控”→“仪表盘”,查看连接数、消息吞吐量、延迟等指标。
步骤3:用MQTTnet连接自己的EMQX
把之前代码中的_brokerAddress改成localhost,_username和_password改成自己创建的用户,运行代码,就能在EMQX控制台看到连接数和消息统计。
步骤4:EMQX集群部署(生产级高可用)
bash
# 启动第一个节点
docker run -d --name emqx1 -p 1883:1883 -p 18083:18083 emqx/emqx:5.6.1
# 获取第一个节点的容器IP
EMQX1_IP=$(docker inspect -f '{{.NetworkSettings.IPAddress}}' emqx1)
# 启动第二个节点,加入集群
docker run -d --name emqx2 -e "EMQX_CLUSTER__DISCOVERY_STRATEGY=static" -e "EMQX_CLUSTER__STATIC__SEEDS=["${EMQX1_IP}:4370"]" emqx/emqx:5.6.1
# 启动第三个节点,加入集群
docker run -d --name emqx3 -e "EMQX_CLUSTER__DISCOVERY_STRATEGY=static" -e "EMQX_CLUSTER__STATIC__SEEDS=["${EMQX1_IP}:4370"]" emqx/emqx:5.6.1
集群节点之间通过4370端口通信;
生产环境建议用负载均衡(比如Nginx)转发MQTT请求到集群节点;
配置持久化(比如Redis、MySQL)存储会话和消息,避免节点挂了丢失数据。
四、生产级最佳实践:让你的MQTT系统稳如老狗
-
客户端优化
用托管客户端:MQTTnet的IManagedMqttClient自动处理重连、消息队列,不用自己写重复代码;
重连机制:设置合理的重连间隔(比如5秒),避免频繁重连导致Broker压力过大;
消息重试:用QoS1或QoS2保证消息可靠到达,托管客户端自动处理重试;
日志记录:记录连接、发布、订阅的日志,方便排查问题(比如用Serilog、NLog);
资源释放:程序退出时调用StopAsync断开连接,避免资源泄漏。 -
Broker优化
集群部署:用EMQX集群保证高可用,避免单点故障;
ACL权限:严格配置主题权限,比如设备只能订阅自己的控制主题,发布自己的状态主题;
TLS加密:生产环境必须启用TLS(端口8883),避免消息被窃听或篡改;
规则引擎:用EMQX规则引擎转发消息到数据库、Kafka、HTTP服务,减少客户端逻辑;
监控告警:配置Prometheus+Grafana监控EMQX,设置连接数、消息延迟的告警阈值。 -
安全配置
禁用匿名访问:EMQX默认允许匿名连接,生产环境必须关闭,启用用户名密码认证;
证书认证:用X.509证书认证设备,比用户名密码更安全;
限流限速:限制单个客户端的消息发送频率,避免恶意攻击;
数据加密:敏感消息(比如用户隐私)在客户端加密后再发送,Broker只存储密文。
五、总结与选型建议 -
总结
MQTTnet是C# MQTT开发的首选库,异步高性能,功能丰富,跨平台;
EMQX是生产级MQTT Broker的最佳选择,百万级并发,分布式集群,全链路监控;
生产环境必须结合托管客户端+EMQX集群+安全配置,才能保证系统的高可用和可靠性。 - 选型建议
| 场景 | 技术栈 |
|---|---|
| 智能家居、车联网 | MQTTnet客户端 + EMQX集群 + TLS加密 |
| 工业物联网 | MQTTnet客户端 + EMQX集群 + 规则引擎转发到工业数据库 |
| 低功耗设备 | MQTTnet Lite(针对受限设备优化) + EMQX边缘节点 |
| 云原生物联网 | MQTTnet客户端 + EMQX Cloud(托管服务) + Kubernetes部署 |
下一节我们会学习C# CoAP编程,用CoAP.NET库开发低功耗物联网节点,结合EMQX的CoAP网关实现跨协议通信!
转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49578.html










