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

多并发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()

问题根源

  1. 重复Client ID引发连接冲突:全局定义的client_id仅生成一次,循环创建客户端时所有实例共用同一个ID。MQTT Broker不允许同一Client ID存在多个活跃连接,后创建的客户端会踢掉前一个连接,导致反复断开重连的循环。
  2. 冗余的多客户端设计:一个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 22:52:31