You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用C# M2Mqtt库MQTT消息丢失:QOS2仍无法接收离线消息

问题分析与解决方案

核心原因

离线消息丢失和QoS等级本身无关,关键是MQTT离线消息存储的前提条件未满足。MQTT Broker仅会为符合以下条件的客户端存储离线消息:

  • 客户端连接时设置cleanSession=false(持久化会话),且客户端ID固定唯一
  • 客户端订阅的QoS等级 ≥ 消息发送的QoS等级
  • Broker启用了持久化存储(默认部分Broker未配置)

你的代码存在多处关键问题,导致离线消息无法被存储和投递:

代码问题与修复

1. 接收端:连接顺序错误+未启用持久化会话

接收端先调用Subscribe再Connect,会导致订阅请求无法被Broker持久化;同时默认cleanSession=true,Broker不会保存该客户端的订阅状态和离线消息。

修复后的接收端代码:

public class Programm
{
    static MqttClient mqttClient;

    static async Task Main(string[] args)
    {
        var clientName = "Emfänger 1";
        var locahlost = true;

        Console.WriteLine($"Start of {clientName}");

        Task.Run(() =>
        {
            var servr = locahlost ? "localhost" : "test.mosquitto.org";
            mqttClient = new MqttClient(servr);
            mqttClient.MqttMsgPublishReceived += MqttClient_MqttMsgPublishReceived;
            
            // 关键:设置cleanSession=false,启用持久化会话
            mqttClient.Connect(clientName, null, null, false, MqttMsgBase.QOS_LEVEL_AT_MOST_ONCE);
            
            // 连接成功后再订阅,确保订阅被Broker持久化
            mqttClient.Subscribe(new string[] { "Application1/NEW_Message" }, new byte[] { MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE });
        });

        Console.ReadLine();
        Console.WriteLine($"end of  {clientName}");
        Console.ReadLine();
    }

    private static void MqttClient_MqttMsgPublishReceived(object sender, uPLibrary.Networking.M2Mqtt.Messages.MqttMsgPublishEventArgs e)
    {
        var message = Encoding.UTF8.GetString(e.Message);
        Console.WriteLine(message);
    }
}

2. 发送端:无需设置retain=true(除非需保留最后一条消息)

你发送消息时设置了retain=true,但这仅会让Broker保留该主题的最后一条消息,并非存储所有离线消息。若无需保留最后一条,建议改为false,避免不必要的存储。

发送端Publish方法修改:

mqttClient.Publish("Application1/NEW_Message", Encoding.UTF8.GetBytes($"{Message}"), MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE, false);

3. 服务器端:启用持久化存储

你使用的MQTTnet服务器默认是内存存储,重启后会丢失所有订阅和消息。若需服务器重启后仍保留离线消息,需配置文件持久化:

修复后的服务器代码:

public class Programm
{
    static async Task Main(string[] args)
    {
        Console.WriteLine("Server");
        var options = new MqttServerOptionsBuilder()
            .WithDefaultEndpoint()
            .WithDefaultEndpointPort(1883)
            // 启用文件持久化,指定存储路径
            .WithPersistentStorage("mqtt_persistence");

        IMqttServer mqttServer = new MqttFactory().CreateMqttServer();
        await mqttServer.StartAsync(options.Build());

        Console.ReadLine();
        await mqttServer.StopAsync();
    }
}

验证步骤

  1. 启动服务器
  2. 启动接收端,确认连接成功并完成订阅后关闭接收端
  3. 启动发送端发送消息
  4. 再次启动接收端,此时应能收到离线期间发送的所有消息

内容的提问来源于stack exchange,提问作者M.Michael

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 21:01:04