MQTT订阅端无法主动接收消息、持久会话配置不生效问题咨询
MQTT订阅端无法实时接收消息问题排查与解决
核心问题根因
- 客户端ID冲突:发布端和订阅端使用了完全相同的clientId
clientId,MQTT协议规定同一broker下不允许两个相同clientId的客户端同时在线,后连接的客户端会主动踢掉已在线的客户端,导致两端频繁断连,无法正常收发消息。 - 订阅逻辑顺序错误:订阅端代码中先执行
Subscribe订阅请求,再执行Connect连接broker,未建立连接时发送的订阅请求不会被broker受理,仅首次启动时的时序巧合可能让订阅临时生效,因此只能收到一次消息。 - 缺少断连自动重连逻辑:现有代码没有监听连接断开事件,也没有重连机制,一旦网络波动或被踢下线,订阅端不会自动恢复连接,自然无法接收后续消息。
修复方案
1. 修正客户端ID配置
给发布端和订阅端分配完全不同的固定clientId,保证持久会话生效的同时避免冲突。
2. 调整订阅端执行顺序
先建立broker连接,连接成功后再发送订阅请求。
3. 新增断连重连逻辑
监听连接断开事件,触发后自动重试连接,重连成功后可按需重新发起订阅。
修正后代码示例
发布端代码
public class publisher : IDatabaseSubscription { private MqttClient _mqttClient; // 发布端固定唯一clientId private const string PublisherClientId = "publisher_client_001"; public publisher() { _mqttClient = new MqttClient("127.0.0.1"); _mqttClient.Connect(PublisherClientId, null, null, false, 60); // 添加断连重连逻辑 _mqttClient.ConnectionClosed += (s, e) => { Thread.Sleep(3000); if (!_mqttClient.IsConnected) { _mqttClient.Connect(PublisherClientId, null, null, false, 60); } }; } private void MQTT_OnChanged(object sender, RecordChangedEventArgs<EventDto> e) { if (_mqttClient != null && _mqttClient.IsConnected) { var message = System.Text.Encoding.UTF8.GetBytes("Hello from app 1"); var statusCode = _mqttClient.Publish("Message1",message , MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE, true); } } }
订阅端代码
public class subscriber { private MqttClient _mqttClient; // 订阅端固定唯一clientId,持久会话依赖该ID保持不变 private const string SubscriberClientId = "subscriber_client_001"; public subscriber() { } public void configure () { _mqttClient = new MqttClient("127.0.0.1"); // 先注册消息接收回调 _mqttClient.MqttMsgPublishReceived += client_MqttMsgPublishReceived; // 先连接broker _mqttClient.Connect(SubscriberClientId, null, null , false, 60) ; // 连接成功后再订阅主题 var status = _mqttClient.Subscribe(new string[] { "Message1" }, new byte[] { MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE }); // 添加断连重连逻辑 _mqttClient.ConnectionClosed += (s, e) => { Thread.Sleep(3000); if (!_mqttClient.IsConnected) { _mqttClient.Connect(SubscriberClientId, null, null, false, 60); // cleanSession为false时理论上不需要重新订阅,可根据broker实际表现选择是否开启 // _mqttClient.Subscribe(new string[] { "Message1" }, new byte[] { MqttMsgBase.QOS_LEVEL_EXACTLY_ONCE }); } }; } private void client_MqttMsgPublishReceived(object sender, MqttMsgPublishEventArgs e) { var message = System.Text.Encoding.UTF8.GetString(e.Message); // 此处添加业务处理逻辑,避免执行长时间阻塞操作 } }
启动项代码
if(IsSubscriptionEnabled()) { // 类名修正为实际定义的subscriber,避免实例化错误 var service = _container.GetInstance<subscriber>(); service.configure(); // 根据服务框架调整逻辑,避免configure执行完后线程退出导致订阅终止 }
额外验证项
- 确认MQTT broker端开启了持久会话支持,部分broker默认会限制持久会话有效期,可调整配置关闭有效期限制。
- 若使用M2Mqtt库,建议升级到最新稳定版,避免已知的订阅逻辑bug。
- 消息接收回调中不要执行耗时过长的阻塞逻辑,避免客户端消息队列阻塞导致后续消息无法接收。
内容的提问来源于stack exchange,提问作者luci
相关产品推荐
相关产品推荐

