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

Mosquitto接收QoS 2消息后未转发至在线已订阅客户端的问题排查求助

Mosquitto接收QoS 2消息后未转发至在线已订阅客户端的问题排查求助

我已经被这个MQTT问题折腾快一周了,问题看起来特别隐蔽,排查起来异常棘手。先给大家说下我们的场景:我们有大概70台嵌入式设备,通过Mosquitto broker和后端服务进行传感器数据传输,整个消息流程是这样的:

  • 设备发送一条QoS 2的"start"消息到后端
  • 后端回复一条QoS 2的"ack"消息给设备
  • 设备以QoS 0的方式批量发送传感器数据
  • 后端最终发送一条QoS 2的"stop"消息给设备

这套流程在第一轮执行的时候完全正常,但到了第二轮(或者后续轮次),部分设备的"start"消息就再也传不到后端了——更诡异的是,我已经确认了以下几点:

  • 设备确实发送了"start"消息(Wireshark抓包实锤)
  • Mosquitto broker已经接收到这条消息(日志和Wireshark都能看到)
  • 后端客户端确实订阅了正确的主题(日志和客户端状态都验证过)
  • 后端客户端全程保持连接且订阅状态没变化

更头疼的是这个问题不是必现的,完全随机挑设备出问题。我现在满脑子都是问号:为什么Mosquitto收到了QoS 2的消息,却不转发给一个在线且已订阅的客户端? 尤其是在满足这些条件的情况下:

  • 消息的主题和payload和之前完全一致
  • 每次的Message ID都是不同的
  • 客户端全程保持连接且订阅状态没变化

有没有大佬能给我分析下可能的原因?或者推荐一些能进一步排查的工具?真的万分感谢🙏

以下是故障设备的Wireshark抓包截图:
[故障设备抓包截图]

对了,我的MQTT客户端初始化和连接代码是这样的:

客户端初始化代码

def __init_client_tx(self) -> bool:
    is_initialized = False
    try:
        client_id = self.cfg_mqtt.get("client_id", "")
        if not isinstance(client_id, str):
            raise TypeError("client_id must be a string")
        api_version = getattr(paho.CallbackAPIVersion, "VERSION2", None)
        valid_protocols = {paho.MQTTv31, paho.MQTTv311, paho.MQTTv5}
        protocol = self.cfg_mqtt.get("protocol", paho.MQTTv5)
        if protocol not in valid_protocols:
            raise ValueError(f"Invalid protocol: {protocol}")
        self.__client_tx = paho.Client(
            callback_api_version=api_version,
            client_id=client_id + "_tx",
            protocol=protocol,
        )
        is_initialized = True
    except (KeyError, AttributeError, TypeError, ValueError, Exception) as e:
        pass
    return is_initialized

def __init_client_rx(self) -> bool:
    is_initialized = False
    try:
        client_id = self.cfg_mqtt.get("client_id", "")
        if not isinstance(client_id, str):
            raise TypeError("client_id must be a string")
        api_version = getattr(paho.CallbackAPIVersion, "VERSION2", None)
        valid_protocols = {paho.MQTTv31, paho.MQTTv311, paho.MQTTv5}
        protocol = self.cfg_mqtt.get("protocol", paho.MQTTv5)
        if protocol not in valid_protocols:
            raise ValueError(f"Invalid protocol: {protocol}")
        self.__client_rx = paho.Client(
            callback_api_version=api_version,
            client_id=client_id + "_rx",
            protocol=protocol,
        )
        is_initialized = True
    except (KeyError, AttributeError, TypeError, ValueError, Exception) as e:
        pass
    return is_initialized

客户端连接代码

def __connect_tx(self) -> MQTTErrorCode:
    connect_res: MQTTErrorCode = MQTTErrorCode.MQTT_ERR_UNKNOWN
    try:
        self.__client_tx.username_pw_set(
            username=os.environ[self.cfg_mqtt["userAccess"]],
            password=os.environ[self.cfg_mqtt["passwordAccess"]],
        )
        # 绑定MQTT回调函数
        self.__client_tx.on_connect = self.on_connect_tx
        self.__client_tx.on_message = self.on_message_tx
        self.__client_tx.on_disconnect = self.on_disconnect_tx
        # 连接对应环境的MQTT服务器
        valid_envs = {"LOCAL", "DEVELOPMENT", "PREPRODUCTION", "PRODUCTION"}
        if self.env not in valid_envs:
            sys.exit("ERROR: Invalid environment! Choose from LOCAL, DEVELOPMENT, PREPRODUCTION, or PRODUCTION.")
        else:
            mqtt_host = self.cfg_mqtt["address"][self.env.lower()]
            mqtt_port = self.cfg_mqtt["port"]
            # 创建CONNECT报文的属性对象
            connect_properties = paho.Properties(paho.PacketTypes.CONNECT)
            # 设置MQTT 5相关属性
            connect_properties.SessionExpiryInterval = 600  # 会话10分钟后过期
            connect_properties.ReceiveMaximum = 5  # 限制最大接收的QoS1/QoS2消息数量
            connect_properties.MaximumPacketSize = 65536  # 最大包字节数
            connect_properties.TopicAliasMaximum = 10  # 允许最多10个主题别名
            connect_properties.RequestProblemInformation = 1  # 向broker请求问题详情
            connect_res = self.__client_tx.connect(
                host=mqtt_host,
                port=mqtt_port,
                clean_start=paho.MQTT_CLEAN_START_FIRST_ONLY,
                properties=connect_properties,
            )
    except Exception as e:
        pass
    return connect_res

def __connect_rx(self) -> MQTTErrorCode:
    connect_res: MQTTErrorCode = MQTTErrorCode.MQTT_ERR_UNKNOWN
    try:
        self.__client_rx.username_pw_set(
            username=os.environ[self.cfg_mqtt["userAccess"]],
            password=os.environ[self.cfg_mqtt["passwordAccess"]],
        )
        # 绑定MQTT回调函数
        self.__client_rx.on_connect = self.on_connect_rx
        self.__client_rx.on_message = self.on_message_rx
        self.__client_rx.on_disconnect = self.on_disconnect_rx
        # 连接对应环境的MQTT服务器
        valid_envs = {"LOCAL", "DEVELOPMENT", "PREPRODUCTION", "PRODUCTION"}
        if self.env not in valid_envs:
            sys.exit("ERROR: Invalid environment! Choose from LOCAL, DEVELOPMENT, PREPRODUCTION, or PRODUCTION.")
        else:
            mqtt_host = self.cfg_mqtt["address"][self.env.lower()]
            mqtt_port = self.cfg_mqtt["port"]
            # 创建CONNECT报文的属性对象
            connect_properties = paho.Properties(paho.PacketTypes.CONNECT)
            # 设置MQTT 5相关属性
            connect_properties.SessionExpiryInterval = 600  # 会话10分钟后过期
            connect_properties.ReceiveMaximum = 5  # 限制最大接收的QoS1/QoS2消息数量
            connect_properties.MaximumPacketSize = 65536  # 最大包字节数
            connect_properties.TopicAliasMaximum = 10  # 允许最多10个主题别名
            connect_properties.RequestProblemInformation = 1  # 向broker请求问题详情
            connect_res = self.__client_rx.connect(
                host=mqtt_host,
                port=mqtt_port,
                clean_start=paho.MQTT_CLEAN_START_FIRST_ONLY,
                properties=connect_properties,
            )
    except Exception as e:
        pass
    return connect_res

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 09:58:01