MQTT - Snet Docs

📨 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)。MqttServiceOperateMqttWebSocketServiceOperate 是服务端类,直接继承 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 的异步变体

另请参阅