多并发MQTT订阅者无法持续在线问题排查与实现方案咨询
MQTT多订阅者反复重连问题排查与优化方案
问题说明
尝试实现多发布者-多订阅者的MQTT消息交互,启动多个订阅不同主题的订阅者时出现反复重连现象,期望订阅端能持续在线等待发布端消息,当前订阅端代码如下:
import random from paho.mqtt import client as mqtt_client import threading broker='172.18.0.3' port = 1885 topics = ["sensors/humidity","sensors/temperature"] # Generate a Client ID with the subscribe prefix. client_id = f'subscribe-{random.randint(0, 100)}' # ++ OBJECT ORIENTED PROGRAMMING? def connect_mqtt() -> mqtt_client: def on_connect(client, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker!") else: print("Failed to connect, return code %d\n", rc) client = mqtt_client.Client(client_id) # client.username_pw_set(username, password) client.on_connect = on_connect client.connect(broker, port) return client def subscribe(client: mqtt_client, topic): def on_message(client, userdata, msg): print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic") client.subscribe(topic) client.on_message = on_message def run(): threads=[] for topic in topics: client = connect_mqtt() subscribe(client, topic) thread=threading.Thread(target=client.loop_forever) threads.append(thread) for thread in threads: thread.start() if __name__ == '__main__': run()
问题根源
- 重复Client ID引发连接冲突:全局定义的
client_id仅生成一次,循环创建客户端时所有实例共用同一个ID。MQTT Broker不允许同一Client ID存在多个活跃连接,后创建的客户端会踢掉前一个连接,导致反复断开重连的循环。 - 冗余的多客户端设计:一个MQTT客户端可同时订阅多个主题,无需为每个主题单独创建客户端和线程,这种设计既浪费资源又容易引发连接问题。
优化实现方案
方案1:单客户端多主题订阅(推荐)
这是最简洁高效的方式,一个客户端即可完成多主题订阅,无需多线程:
import random from paho.mqtt import client as mqtt_client broker = '172.18.0.3' port = 1885 topics = ["sensors/humidity", "sensors/temperature"] # 生成唯一Client ID client_id = f'subscribe-{random.randint(0, 1000)}' def connect_mqtt() -> mqtt_client: def on_connect(client, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker!") # 连接成功后批量订阅所有主题 for topic in topics: client.subscribe(topic) print(f"Subscribed to {topic}") else: print(f"Failed to connect, return code {rc}") def on_message(client, userdata, msg): print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic") client = mqtt_client.Client(client_id) client.on_connect = on_connect client.on_message = on_message client.connect(broker, port) return client def run(): client = connect_mqtt() client.loop_forever() if __name__ == '__main__': run()
方案2:多客户端(需保证Client ID唯一)
如果业务场景确实需要多客户端,必须确保每个客户端使用唯一的Client ID:
import random from paho.mqtt import client as mqtt_client import threading broker = '172.18.0.3' port = 1885 topics = ["sensors/humidity", "sensors/temperature"] def connect_mqtt(client_id) -> mqtt_client: def on_connect(client, userdata, flags, rc): if rc == 0: print(f"Client {client_id} connected to MQTT Broker!") else: print(f"Client {client_id} failed to connect, return code {rc}") client = mqtt_client.Client(client_id) client.on_connect = on_connect client.connect(broker, port) return client def subscribe(client: mqtt_client, topic): def on_message(client, userdata, msg): print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic (Client: {client._client_id.decode()})") client.subscribe(topic) client.on_message = on_message def run(): threads = [] for idx, topic in enumerate(topics): # 为每个客户端生成唯一ID client_id = f'subscribe-{random.randint(0, 1000)}-{idx}' client = connect_mqtt(client_id) subscribe(client, topic) thread = threading.Thread(target=client.loop_forever) threads.append(thread) for thread in threads: thread.start() if __name__ == '__main__': run()
补充说明
- MQTT 3.1/3.1.1协议中,同一Client ID的新连接会强制断开旧连接;MQTT 5.0可通过配置允许共享订阅,但默认行为仍为踢掉旧连接。
- 单客户端多主题订阅是MQTT的标准用法,能有效减少Broker连接负载,简化客户端维护。
内容的提问来源于stack exchange,提问作者Mattia Sospetti
相关产品推荐
相关产品推荐

