如何在Paho MQTT的on_connect方法中获取动态主题提取设备序列号?
MQTT客户端提取设备序列号的实现方案
核心结论
不能在on_connect方法中获取代理发布的主题。on_connect仅在客户端与代理建立连接完成时触发,此时客户端还未订阅任何主题,也不会收到代理主动推送的消息,自然无法获取到device/****/report这类主题内容。
正确实现步骤
- 在连接成功的
on_connect回调里,订阅带通配符的主题device/+/report(MQTT的+单级通配符可匹配任意单个层级的动态值,正好对应设备序列号)。 - 实现
on_message回调,当收到消息时,从消息的topic字段中拆分出设备序列号。
修改后的代码示例
def try_connection( user_input: dict[str, Any], ) -> bool: """Test if we can connect to an MQTT broker and extract device serial number.""" # We don't import on the top because some integrations # should be able to optionally rely on MQTT. import paho.mqtt.client as mqtt # pylint: disable=import-outside-toplevel import queue client = mqtt.Client() result: queue.Queue[bool] = queue.Queue(maxsize=1) serial_number_queue: queue.Queue[str] = queue.Queue(maxsize=1) def on_connect( client_: mqtt.Client, userdata: None, flags: dict[str, Any], result_code: int, properties: mqtt.Properties | None = None, ) -> None: """Handle connection result and subscribe to target topic.""" LOGGER.debug(f"client: {client.__dict__}") LOGGER.debug(f"flags: {flags}") LOGGER.debug(f"result_code: {result_code}") LOGGER.debug(f"properties: {properties}") if result_code == mqtt.CONNACK_ACCEPTED: # 订阅带通配符的主题,匹配所有device/{序列号}/report格式的主题 client_.subscribe("device/+/report") result.put(True) else: result.put(False) def on_message( client_: mqtt.Client, userdata: None, msg: mqtt.MQTTMessage, ) -> None: """Handle received messages and extract serial number from topic.""" # 拆分主题,格式为device/{serial}/report,索引1的位置就是序列号 topic_parts = msg.topic.split("/") if len(topic_parts) == 3 and topic_parts[0] == "device" and topic_parts[2] == "report": serial_number = topic_parts[1] LOGGER.debug(f"Extracted device serial number: {serial_number}") serial_number_queue.put(serial_number) client.on_connect = on_connect client.on_message = on_message client.connect_async(user_input[CONF_HOST], 1883) client.loop_start() try: # 先确认连接成功 if result.get(timeout=5): # 尝试获取解析到的序列号,可根据实际需求调整超时时间 try: serial = serial_number_queue.get(timeout=10) LOGGER.info(f"Successfully got device serial number: {serial}") except queue.Empty: LOGGER.warning("No message received within timeout period") return True return False except queue.Empty: return False finally: client.disconnect() client.loop_stop()
内容的提问来源于stack exchange,提问作者K20GH
相关产品推荐
相关产品推荐

