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

使用M2Mqtt.net对接Mosquitto 2.0.12时Retained MQTT消息丢失问题

问题根因及修复方案

核心错误点

  • 客户端ID使用动态GUID:持久会话(cleanSession=false)的核心是固定客户端ID,broker通过客户端ID关联离线消息队列,每次生成新的GUID相当于每次都是全新的客户端,broker不会为不存在的客户端缓存离线消息。
  • 订阅操作顺序错误:先调用Subscribe再执行Connect的逻辑不符合M2Mqtt库的要求,未连接状态下调用订阅接口不会生效,订阅请求根本没有发送到broker。
  • Retain标志误用:发布时将retain设为true,Retain机制只会存储主题最后一条消息,和离线消息队列完全是两个功能,收全量离线消息不需要设置retain为true。
  • 接收端未提前创建持久会话:持久会话需要接收端至少成功连接一次broker,broker才会为这个客户端ID创建会话缓存,先发布消息再首次连接接收端,之前发布的消息broker根本不会缓存。
  • 编码不统一:发布时用Encoding.UTF8转码消息,接收时用Encoding.Default,特殊场景下会出现消息解析错误,导致匹配不到Total里的key,误认为丢消息。

修复步骤

  1. 给接收端设置固定的客户端ID,不要用动态GUID
  2. 调整订阅逻辑:连接成功后再执行订阅操作
  3. 测试离线消息场景时,先让接收端连接一次broker创建持久会话,再断开后执行发布,最后重连接收端收离线消息
  4. 发布时不需要设置retain为true,关闭retain避免干扰
  5. 收发两端统一使用UTF8编码解析消息

修正后的核心代码示例

接收端代码

public class MessageReceiver : IMessageReceiver
{
    private readonly MqttClient _client;
    // 固定客户端ID,持久会话绑定用
    private const string ReceiverClientId = "Fixed_Receiver_Client_001";

    public MessageReceiver()
    {
        _client = new MqttClient("localhost");
        _client.MqttMsgPublishReceived += client_receivedMessage;
    }

    public void Subscribe(params string[] topics)
    {
        // 订阅必须在连接成功后调用
        _client.Subscribe(topics, new[] { MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE });
    }

    public void Connect()
    {
        // 使用固定客户端ID,cleanSession设为false
        _client.Connect(ReceiverClientId, "username", "password", false, byte.MaxValue);
    }

    public void Disconnect()
    {
        _client.Disconnect();
    }

    static void client_receivedMessage(object sender, MqttMsgPublishEventArgs e)
    {
        // 统一用UTF8编码
        var message = Encoding.UTF8.GetString(e.Message);
        Console.WriteLine($"Message Received: {message}");
        if (Total.SentAndReceived.ContainsKey(message))
            Total.SentAndReceived[message] = message;
    }
}

发布端代码

public class MessagePublisher : IMessagePublisher
{
    private readonly MqttClient _client;
    private const string PublisherClientId = "Fixed_Publisher_Client_001";

    public MessagePublisher()
    {
        _client = new MqttClient("localhost");
        _client.Connect(PublisherClientId, "username", "password", false, byte.MaxValue);
    }

    public void Publish(string topic, string message, bool retain = false)
    {
        Console.Write($"Sent: {topic}, {message}");
        // 不需要开启retain
        _client.Publish(topic, Encoding.UTF8.GetBytes(message), MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE, retain);
        Total.SentAndReceived.Add(message, null);
    }
}

测试主程序逻辑(模拟离线消息场景)

public static class Program
{
    public static void Main(string[] args)
    {
        var messageReceiver = new MessageReceiver();
        // 第一步:先让接收端连接一次,创建持久会话,然后断开
        messageReceiver.Connect();
        messageReceiver.Subscribe("Devices/");
        messageReceiver.Disconnect();

        var publisher = new MessagePublisher();
        // 第二步:接收端离线状态下发布消息
        for (var i = 1; i <= 10000; i++)
        {
            // 关闭retain
            publisher.Publish("Devices/", i.ToString(), false);
        }

        // 第三步:接收端重连,拉取离线消息
        messageReceiver.Connect();

        Thread.Sleep(10000); // 延长等待时间保证10000条消息全部处理完

        var b = Total.SentAndReceived.Where(x => x.Value == null);
        Console.WriteLine($"{b.Count()} Missed Messages");
    }
}

额外注意点

  • mosquitto默认的持久会话消息存储有大小和数量限制,如果消息量太大需要调整mosquitto配置里的max_queued_messages、message_size_limit参数,避免队列满了丢弃旧消息(这也是丢失前n条消息的核心原因之一)
  • QoS2的消息处理需要等待确认,发布端如果发的太快可能会导致本地发送队列溢出,可以适当加发送间隔,或者调整库的发送缓冲区配置

内容的提问来源于stack exchange,提问作者Kieran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 20:54:02