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

MQTT订阅端无法主动接收消息、持久会话配置不生效问题咨询

MQTT订阅端无法实时接收消息问题排查与解决

核心问题根因

  • 客户端ID冲突:发布端和订阅端使用了完全相同的clientIdclientId,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 08:24:05