之前项目里用 MQTTnet 做 MQTT 通信,一开始用的是普通的 MqttClient,断线了得自己写重连逻辑,写得又臭又长。后来发现 MQTTnet 有个扩展包叫 ManagedMqttClient,断线重连、离线消息排队这些它都帮你干了,用起来舒服很多。

另外 MQTT 里的 QoS、CleanSession、Retain 这几个东西我之前也一直记不清谁是谁,正好一起记一下。


安装

dotnet add package MQTTnet
dotnet add package MQTTnet.Extensions.ManagedClient

ManagedMqttClient 是一个独立的扩展包,不在 MQTTnet 主包里。

基本用法

var options = new ManagedMqttClientOptionsBuilder()
    .WithAutoReconnectDelay(TimeSpan.FromSeconds(5))
    .WithClientOptions(new MqttClientOptionsBuilder()
        .WithClientId("my-client-001")
        .WithTcpServer("broker.hivemq.com", 1883)
        .WithCredentials("username", "password")
        .WithCleanSession(false)
        .Build())
    .Build();

var mqttClient = new MqttFactory().CreateManagedMqttClient();

// 收到消息的回调
mqttClient.ApplicationMessageReceivedAsync += e =>
{
    var topic = e.ApplicationMessage.Topic;
    var payload = Encoding.UTF8.GetString(e.ApplicationMessage.PayloadSegment);
    Console.WriteLine($"收到消息 [{topic}]: {payload}");
    return Task.CompletedTask;
};

// 连接状态变化
mqttClient.ConnectedAsync += e =>
{
    Console.WriteLine("已连接到 Broker");
    return Task.CompletedTask;
};

mqttClient.DisconnectedAsync += e =>
{
    Console.WriteLine("连接断开,等待重连...");
    return Task.CompletedTask;
};

// 订阅主题
await mqttClient.SubscribeAsync(
    new MqttTopicFilterBuilder()
        .WithTopic("device/sensor/temperature")
        .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce)
        .Build());

// 启动
await mqttClient.StartAsync(options);

StartAsync 调完就返回了,它内部自己开了线程去维护连接。断线之后会按照 WithAutoReconnectDelay 设的间隔自动重连,重连成功后之前的订阅也会自动恢复,不用你再手动 Subscribe 一遍。

发布消息

await mqttClient.EnqueueAsync(
    new MqttApplicationMessageBuilder()
        .WithTopic("device/sensor/temperature")
        .WithPayload("26.5")
        .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce)
        .WithRetainFlag(false)
        .Build());

注意这里用的是 EnqueueAsync 不是 PublishAsync。消息会先进入内部队列,连接正常的时候自动发出去,如果当前是断开状态,消息会在队列里等着,重连成功后再发。这是 ManagedMqttClient 比普通 MqttClient 好用的地方之一。

QoS(服务质量等级)

MQTT 有三个 QoS 等级,控制消息的投递保证:

QoS 0 — 最多一次(At Most Once)

发了就不管了,不确认,不重试。最快,但可能丢消息。适合那种丢了也无所谓的场景,比如传感器每秒上报一次温度,丢一两条没影响。

QoS 1 — 至少一次(At Least Once)

Broker 收到消息后会回一个确认(PUBACK),如果发送方没收到确认就会重发。保证消息不丢,但可能收到重复的。大多数场景用这个就够了。

QoS 2 — 恰好一次(Exactly Once)

通过四次握手(PUBLISH → PUBREC → PUBREL → PUBCOMP)保证消息不丢也不重复。最可靠,但也最慢,开销最大。支付通知、指令下发这种不能重复的场景才需要用。

实际项目里大部分用 QoS 1,QoS 2 能不用就不用,性能差别还是挺明显的。

CleanSession(清除会话)

连接 Broker 的时候有个 CleanSession 开关:

CleanSession = true

每次连接都是全新的,Broker 不会保留任何和你相关的状态。断线期间别人发给你的消息全部丢弃,重连后需要重新订阅。

CleanSession = false

Broker 会记住你的订阅关系和你离线期间收到的 QoS 1/2 消息。你重连上来之后,Broker 会把这些积压的消息推给你。前提是你的 ClientId 要固定,不能每次连接都随机生成。

如果你的场景是:设备偶尔断网,重连之后要把断网期间的消息补回来,就把 CleanSession 设成 false。如果不在意离线消息,设成 true 就行,Broker 的压力也小一些。

Retain(保留消息)

发布消息的时候可以设一个 Retain 标志:

Retain = true

Broker 会保留这个 Topic 上最后一条带 Retain 标志的消息。新的订阅者订阅这个 Topic 的时候,会立刻收到这条保留消息,不用等下一次发布。

Retain = false

消息发完就完了,新订阅者订阅之后要等下一次发布才能收到消息。

一个典型的场景:设备上线后发一条 online 状态消息并设置 Retain,这样任何时候有新的客户端订阅这个设备的状态 Topic,都能马上知道它是在线的,不用等设备下一次上报。

要清除某个 Topic 的保留消息,发一条 payload 为空、Retain 为 true 的消息就行。

这几个东西怎么配合

举个例子:一个温度传感器每 10 秒上报一次数据,监控面板订阅这个数据。

传感器端发布:QoS 1 + Retain = true。QoS 1 保证消息不丢,Retain 保证监控面板刷新页面后马上能看到最新的温度值,不用等 10 秒。

监控面板订阅:QoS 1 + CleanSession = true。面板不需要补历线消息,每次打开看到当前值就够了,Retain 那条消息解决了"打开就能看到"的问题。

如果换成是报警系统,那就不一样了:QoS 1 或 2 + CleanSession = false,断线期间的报警消息一条都不能丢,重连后要全部补回来。