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

如何在Python中分别订阅两个MQTT主题?问题排查及优化建议

MQTT多主题订阅问题优化方案

问题根源

你的代码里client.loop_forever()是阻塞式调用,会一直占用当前线程执行MQTT消息循环,第一次调用subscribe_topic后主线程被卡住,永远无法执行第二个客户端的创建和连接逻辑。另外全局变量connected和MessageReceived会被多个客户端共享,导致状态混乱,也是潜在问题。


方案一:单客户端订阅多个主题(推荐)

没必要为每个主题单独创建客户端,一个MQTT客户端可以同时订阅多个主题,既节省资源又避免线程阻塞问题。

优化后代码

import paho.mqtt.client as mqttclient
from excel_adapter import excel_yaz

def on_connect(client, userdata, flags, rc):
    if rc == 0:
        print(f"{client._client_id.decode('utf-8')} 已连接")
        # 订阅所有传入的主题
        for topic in userdata['topics']:
            client.subscribe(topic)
            print(f"已订阅主题: {topic}")
    else:
        print(f"连接失败,错误码: {rc}")

def on_message(client, userdata, msg):
    payload = str(msg.payload.decode('utf-8'))
    print(f"{client._client_id.decode('utf-8')}: {msg.topic} {payload}")
    excel_yaz(payload)

def start_mqtt_client(client_name, topics):
    mqtt_port = 1883
    mqtt_broker = "xxxxxx"
    mqtt_username = "yyyyyy"
    mqtt_password = "111111"

    # 用userdata传递主题列表,避免全局变量
    client = mqttclient.Client(client_name)
    client.username_pw_set(mqtt_username, password=mqtt_password)
    client.user_data_set({'topics': topics})
    client.on_connect = on_connect
    client.on_message = on_message
    
    client.connect(mqtt_broker, port=mqtt_port)
    # 启动后台消息循环,不阻塞主线程
    client.loop_start()

调用方式(main.py)

from your_module import start_mqtt_client

# 传入要订阅的主题列表
start_mqtt_client("Client1", ["topic/1", "topic/2"])
# 如果需要多个独立客户端(比如不同权限),直接再调用一次
start_mqtt_client("Client2", ["topic/3", "topic/4"])

# 主线程保持运行
while True:
    pass

方案二:多线程处理多客户端(适用于必须用多个客户端的场景)

如果确实需要多个独立客户端,要把每个客户端的消息循环放到单独线程中执行,避免阻塞主线程。

优化后代码

import paho.mqtt.client as mqttclient
import threading
from excel_adapter import excel_yaz

def subscribe_topic(topic, client_name):
    def on_connect(client, userdata, flags, rc):
        if rc == 0:
            print(f"{client_name} 已连接")
            client.subscribe(topic)
        else:
            print(f"{client_name} 连接失败")

    def on_message(client, userdata, msg):
        payload = str(msg.payload.decode('utf-8'))
        print(f"{client_name}: {msg.topic} {payload}")
        excel_yaz(payload)

    mqtt_port = 1883
    mqtt_broker = "xxxxxx"
    mqtt_username = "yyyyyy"
    mqtt_password = "111111"

    client = mqttclient.Client(client_name)
    client.username_pw_set(mqtt_username, password=mqtt_password)
    client.on_connect = on_connect
    client.on_message = on_message
    client.connect(mqtt_broker, port=mqtt_port)
    
    # 启动后台线程执行消息循环,不阻塞当前函数
    client.loop_start()

调用方式(main.py)

from your_module import subscribe_topic
import threading

# 用线程分别启动两个客户端
threading.Thread(target=subscribe_topic, args=("topic/1", "Client1")).start()
threading.Thread(target=subscribe_topic, args=("topic/2", "Client2")).start()

# 主线程保持运行
while True:
    pass

关键优化点说明

  1. 替换loop_forever()为loop_start():loop_start()会在后台启动一个线程处理MQTT消息循环,不会阻塞主线程,多个客户端可以同时运行。
  2. 移除全局变量:改用userdata传递自定义数据(如主题列表),避免多个客户端共享状态导致的逻辑混乱。
  3. 单客户端多主题订阅:这是MQTT的标准用法,高效且节省资源,优先推荐这种方式。

内容的提问来源于stack exchange,提问作者MrElectronicEngineer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 22:47:49