MQTT Broker未按时投递发布端消息,求问题原因及代码修改方案
问题根因
- 核心问题是paho.mqtt客户端的事件循环被阻塞:
loop_forever()会在当前线程持续处理网络收发、回调触发等核心逻辑,所有的回调函数(包括on_connect)都在这个事件循环线程中执行。 - 你在
on_connect回调中直接调用sendMsg(),而sendMsg()里的while循环加time.sleep(1)会把事件循环线程阻塞5秒。这段时间内客户端无法将本地缓冲区的消息实际发送到Broker,也无法处理Broker返回的发布确认,因此所有消息都被暂存,阻塞结束后才一次性批量发出,就出现了日志批量打印、订阅端批量接收的现象。
修复方案
推荐将发送逻辑和事件循环线程拆分,避免阻塞:
修改后的发布者代码(新增独立线程处理发送)
import paho.mqtt.client as mqtt import time import threading HOST = "localhost" PORT = 1883 KEEP_ALIVE_INT = 100 TOPIC = "noti" def sendMsg(): MSG = ["1111", "2222", "3333", "4444", "5555"] i = 0 try: while i<5: client.publish(TOPIC, MSG[i], qos=0) i+=1 time.sleep(1) except Exception as e: print("Caught Exception: " + str(e)) def onConnect(client, userdata, flags, rc): if rc == 0: print("Connected successfully") # 开独立线程跑发送逻辑,不阻塞事件循环 threading.Thread(target=sendMsg, daemon=True).start() else: print("Connection failed, result code: " + str(rc)) def onPublish(client, userdata, mid): print ("Message is published") client = mqtt.Client("pub") client.on_connect = onConnect client.on_publish = onPublish client.connect(HOST, PORT, KEEP_ALIVE_INT) client.loop_forever()
其他可选方案
你也可以用loop_start()开启后台事件循环,直接在主线程处理发送逻辑:
import paho.mqtt.client as mqtt import time HOST = "localhost" PORT = 1883 KEEP_ALIVE_INT = 100 TOPIC = "noti" def onConnect(client, userdata, flags, rc): if rc == 0: print("Connected successfully") else: print("Connection failed, result code: " + str(rc)) def onPublish(client, userdata, mid): print ("Message is published") client = mqtt.Client("pub") client.on_connect = onConnect client.on_publish = onPublish client.connect(HOST, PORT, KEEP_ALIVE_INT) # 开启后台线程跑事件循环 client.loop_start() # 主线程处理发送逻辑 MSG = ["1111", "2222", "3333", "4444", "5555"] i = 0 try: while i<5: client.publish(TOPIC, MSG[i], qos=0) i+=1 time.sleep(1) except Exception as e: print("Caught Exception: " + str(e))
修改后运行即可实现每秒发送一条消息,订阅端也会每秒收到一条消息,onPublish回调也会逐条触发。
内容的提问来源于stack exchange,提问作者keen
相关产品推荐
相关产品推荐

