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
相关产品推荐
相关产品推荐

