使用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,误认为丢消息。
修复步骤
- 给接收端设置固定的客户端ID,不要用动态GUID
- 调整订阅逻辑:连接成功后再执行订阅操作
- 测试离线消息场景时,先让接收端连接一次broker创建持久会话,再断开后执行发布,最后重连接收端收离线消息
- 发布时不需要设置retain为true,关闭retain避免干扰
- 收发两端统一使用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
相关产品推荐
相关产品推荐

