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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 21:54:01