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

如何在Paho MQTT的on_connect方法中获取动态主题提取设备序列号?

MQTT客户端提取设备序列号的实现方案

核心结论

不能在on_connect方法中获取代理发布的主题。on_connect仅在客户端与代理建立连接完成时触发,此时客户端还未订阅任何主题,也不会收到代理主动推送的消息,自然无法获取到device/****/report这类主题内容。

正确实现步骤

  1. 在连接成功的on_connect回调里,订阅带通配符的主题device/+/report(MQTT的+单级通配符可匹配任意单个层级的动态值,正好对应设备序列号)。
  2. 实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 21:30:44