📨 MQTT 中间件
包名: Snet.Mqtt | 类: MqttClientOperate, MqttServiceOperate, MqttWebSocketServiceOperate | 基类: MqAbstract<O, D>(客户端)/ CoreUnify<O, D>(服务端)
提供 MQTT 3.1.1 客户端连接、内置独立 MQTT 代理以及基于 WebSocket 的 MQTT 服务,适用于工业物联网场景中的实时发布/订阅消息。
概述
Snet.Mqtt 提供三种操作模式:
- MqttClientOperate -- 作为客户端连接到外部 MQTT 代理(如 Mosquitto、EMQX、HiveMQ)
- MqttServiceOperate -- 在应用程序内运行嵌入式 MQTT 代理,无需外部代理
- MqttWebSocketServiceOperate -- 为基于浏览器的客户端提供基于 WebSocket 的 MQTT 服务
这三个类的基类各不相同。MqttClientOperate 继承自 MqAbstract<MqttClientOperate, Basics> 并实现 IMq(IProducer + IConsumer)。MqttServiceOperate 和 MqttWebSocketServiceOperate 是服务端类,直接继承 CoreUnify<MqttServiceOperate, Basics>(分别为 CoreUnify<MqttWebSocketServiceOperate, Basics>)并实现 IOn/IOff/IStatus;它们不实现 IMq。
安装
dotnet add package Snet.Mqtt
快速开始
using Snet.Mqtt.client;
var mqttClient = await MqttClientOperate.InstanceAsync(new MqttClientData.Basics
{
IpAddress = "127.0.0.1",
Port = 6688
});
await mqttClient.OnAsync();
// 绑定事件再消费
mqttClient.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"接收到数据: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
await mqttClient.ConsumeAsync("snet/temperature");
// 生产消息(可选)
await mqttClient.ProduceAsync("snet/temperature", "hello mqtt", System.Text.Encoding.UTF8);
// 后续:取消消费
// await mqttClient.UnConsumeAsync("snet/temperature");
// await mqttClient.DisposeAsync();
生产与消费
发布与订阅是 IMq 的一等操作——先绑定数据事件,再消费;发布可用字符串或原始字节:
// 1) 先绑定事件再消费——收到的消息在此到达
mqttClient.OnDataEventAsync += async (sender, e) =>
{
if (e.Status)
Console.WriteLine($"收到: {e.ResultData}");
else
Console.WriteLine($"消费失败: {e.Message}");
};
// 2) 消费:订阅主题
await mqttClient.ConsumeAsync("snet/temperature");
// 3) 生产:发布字符串消息(UTF-8)
await mqttClient.ProduceAsync("snet/temperature", "hello mqtt", System.Text.Encoding.UTF8);
// 4) 生产:发布原始字节
var bytes = System.Text.Encoding.UTF8.GetBytes("payload");
await mqttClient.ProduceAsync("snet/temperature", bytes);
// 后续:停止消费
// await mqttClient.UnConsumeAsync("snet/temperature");
ProduceAsync(topic, string, Encoding?)与ProduceAsync(topic, byte[])均可发布;ConsumeAsync/UnConsumeAsync管理主题订阅。事件见:事件。
MqttClientOperate 配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
SN |
string | — | 序列号 / 设备标识符 |
IpAddress |
string | "127.0.0.1" |
MQTT 代理 IP 地址 |
Port |
int | 6688 | MQTT 代理端口(默认 6688;仅明文 TCP,无 TLS 配置属性) |
UserName |
string | "shunnet" |
认证用户名 |
Password |
string | "shunnet" |
认证密码 |
ClientID |
string | 自动 | 唯一客户端标识符 |
MessageExpirationTime |
int | 86400000 | 消息过期时间(毫秒) |
QualityOfServiceLevel |
enum | AtMostOnce | QoS 级别:AtMostOnce, AtLeastOnce, ExactlyOnce —— 该配置仅作用于遗嘱消息。ProduceAsync 固定 QoS0 + Retain=true;ConsumeAsync 固定 QoS0。控制 QoS/Retain 请用 Publish/AddSubscribe |
ResponseType |
enum | Content | 响应类型 |
MqttServiceOperate 配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
Port |
int | 6688 | 内置代理监听端口 |
MaxNumber |
int | 10000 | 最大并发连接数 |
UserName |
string | "shunnet" |
代理认证用户名 |
Password |
string | "shunnet" |
代理认证密码 |
MqttWebSocketServiceOperate 配置
| 参数 | 类型 | 默认值 | 描述 |
|---|---|---|---|
Port |
int | 6688 | MQTT 代理端口 |
WsPort |
int | 8866 | WebSocket 服务端口 |
Uri |
string | "shun" |
WebSocket 路径段 |
UserName |
string | "shunnet" |
认证用户名 |
Password |
string | "shunnet" |
认证密码 |
支持的操作
| 操作 | 方法 | 描述 |
|---|---|---|
| 连接 / 启动 | OnAsync() |
连接到代理或启动嵌入式服务 |
| 断开 / 停止 | OffAsync() |
断开连接或停止服务 |
| 发布 | ProduceAsync(topic, string, Encoding?) |
发布字符串消息 |
| 发布 | ProduceAsync(topic, byte[]) |
发布原始字节 |
| 订阅 | ConsumeAsync(topic) |
订阅主题 |
| 取消订阅 | UnConsumeAsync(topic) |
取消订阅主题 |
| 发布(全参数) | Publish(topic, content, QoSLevel, Retain) / PublishAsync(topic, content, QoSLevel, Retain, token) |
以显式 QoS 级别与保留标志发布 |
| 订阅(带 QoS) | AddSubscribe(Topic, QoSLevel) / AddSubscribeAsync(Topic, QoSLevel, token) |
以指定 QoS 级别订阅主题 |
| 取消订阅 | RemoveSubscribe(Topic) / RemoveSubscribeAsync(Topic, token) |
取消订阅主题 |
Publish/AddSubscribe/RemoveSubscribe(及异步变体)是控制 QoS/Retain 的唯一途径。ProduceAsync固定 QoS0 + Retain=true;ConsumeAsync固定 QoS0。
自动重连(客户端)
MQTT 客户端在意外断开后自动重连:最多 5 次,退避间隔 5/10/15/20/30 秒,成功触发「重连成功」信息事件;全部失败后调用 OffAsync(true)。KeepAlive 心跳固定为 60 秒。由于会话为干净会话(WithCleanSession),Broker 在断连时会清空主题订阅——重连成功后需重新订阅(ConsumeAsync/AddSubscribe)才能继续收到消息。
QoS 级别
MQTT 支持三种服务质量级别:
| QoS | 名称 | 描述 |
|---|---|---|
| 0 | 最多一次 | 发后即忘,无确认 |
| 1 | 至少一次 | 保证送达,可能重复 |
| 2 | 恰好一次 | 保证送达且无重复 |
事件
| 事件 | 签名 | 描述 |
|---|---|---|
OnDataEvent |
EventHandler<EventDataResult> |
收到订阅主题数据时触发 |
OnDataEventAsync |
EventHandlerAsync<EventDataResult> |
OnDataEvent 的异步变体 |
OnInfoEvent |
EventHandler<EventInfoResult> |
收到信息通知时触发 |
OnInfoEventAsync |
EventHandlerAsync<EventInfoResult> |
OnInfoEvent 的异步变体 |
OnLanguageEvent |
EventHandler<EventLanguageResult> |
语言变更通知时触发 |
OnLanguageEventAsync |
EventHandlerAsync<EventLanguageResult> |
OnLanguageEvent 的异步变体 |
