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

如何基于Paho-MQTT实现固定频率(20Hz)发布QoS 2消息且不阻塞网络循环

如何基于Paho-MQTT实现固定频率(20Hz)发布QoS 2消息且不阻塞网络循环

你遇到的问题核心很明确:time.sleep()阻塞了Paho-MQTT的网络循环线程。原来的代码里用loop_forever()让MQTT的所有网络交互(比如心跳、QoS级别的报文确认)都在当前线程运行,当你在msg_handler里用sleep+循环时,整个线程被卡住,没法及时处理MQTT的底层报文,导致消息只能堆积到sleep结束后批量发送。

下面给你两种可行的解决方案,分别对应“严格固定间隔”和“基于QoS确认后再发下一条”的场景,你可以根据需求选择:

方案一:严格固定50ms间隔发布(适合要求精准频率的场景)

这个方案用后台网络循环+定时器实现,让MQTT的网络处理和定时发布互不阻塞:

修改后的完整代码

import logging
import time
import threading
import paho.mqtt.client as mqtt

# 假设这些是从配置文件读取的参数
broker_address = "你的 broker 地址"
client_id = "你的客户端ID"
begin_client = "begin_topic"
finish_client = "finish_topic"
main_topic = "publish_topic"
msg_size = 1024
msg_amount = 100
msg_freq = 20  # 20Hz,即50ms间隔

class MQTT_Client:
    def on_connect(self, client, userdata, flags, rc):
        if rc == 0:
            logging.info(f"Connected to the broker at {broker_address}")
        else:
            logging.info(f"Error connecting to broker, with code {rc}")
        self.client.subscribe(begin_client, 0)
        self.client.subscribe(finish_client, 0)
        logging.info(f"Subscribed to {begin_client} and {finish_client} topics with QoS 0")

    def on_message(self, client, userdata, msg):
        if str(msg.topic) == begin_client:
            logging.info(f"Start order received from the server")
            self.msg_handler()
        elif str(msg.topic) == finish_client:
            logging.info(f"Finish order received from the server")
            self.client.disconnect()

    def on_publish(self, client, userdata, mid):
        self.counter += 1
        logging.debug(f"Message {mid} published successfully (QoS 2)")

    def on_disconnect(self, client, userdata, rc):
        logging.info(f"Disconnected from broker at {broker_address}")
        # 通知主线程退出
        self.running_event.set()

    def msg_handler(self):
        time.sleep(2)
        self.payload = bytearray(msg_size)
        self.messages_sent = 0
        logging.info(f"Starting publish of {msg_amount} messages with QoS 2")
        # 启动第一个定时发布任务
        self._schedule_next_publish()

    def _schedule_next_publish(self):
        if self.messages_sent < msg_amount:
            # 发布当前消息
            self.client.publish(main_topic, self.payload, qos=2)
            self.messages_sent += 1
            logging.debug(f"Triggered publish for message {self.messages_sent}/{msg_amount}")
            # 调度下一次发布,严格50ms间隔
            threading.Timer(1/msg_freq, self._schedule_next_publish).start()
        else:
            logging.info(f"Publish complete")

    def __init__(self):
        self.counter = 0
        self.messages_sent = 0
        self.payload = None
        self.running_event = threading.Event()  # 用来保持主线程运行
        logging.info(f"Creating MQTT Client with ID {client_id}")
        
        self.client = mqtt.Client(client_id=client_id)
        self.client.on_connect = self.on_connect
        self.client.on_disconnect = self.on_disconnect
        self.client.on_publish = self.on_publish
        self.client.on_message = self.on_message
        
        self.client.connect(broker_address, 1883, 60)
        # 启动后台网络循环,让MQTT处理在单独线程运行,不阻塞定时任务
        self.client.loop_start()
        
        # 让主线程保持运行,直到收到断开连接信号
        self.running_event.wait()

if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    mqtt_client = MQTT_Client()

关键修改点说明

  1. 用loop_start()替代loop_forever():loop_start()会启动一个后台线程专门处理MQTT的网络交互(比如心跳、QoS2的PUBREC/PUBREL报文),主线程可以自由处理定时任务,不会被阻塞。
  2. 用threading.Timer实现定时发布:每次发布一条消息后,调度50ms后的下一次发布,避免了阻塞式的sleep循环。
  3. 用事件保持主线程运行:添加running_event让主线程不会提前退出,直到收到断开连接的信号。

方案二:基于QoS2确认后再发下一条(适合要求严格可靠性的场景)

如果你需要确保上一条QoS2消息完全确认后,再发送下一条(可能频率会略低于20Hz,但可靠性更高),可以修改成在on_publish回调里触发下一次发布:

核心修改部分

def msg_handler(self):
    time.sleep(2)
    self.payload = bytearray(msg_size)
    self.messages_sent = 0
    logging.info(f"Starting publish of {msg_amount} messages with QoS 2")
    # 发送第一条消息
    self.client.publish(main_topic, self.payload, qos=2)
    self.messages_sent += 1

def on_publish(self, client, userdata, mid):
    self.counter += 1
    logging.debug(f"Message {mid} confirmed (QoS 2)")
    # 确认后发送下一条,直到达到指定数量
    if self.messages_sent < msg_amount:
        self.client.publish(main_topic, self.payload, qos=2)
        self.messages_sent += 1
        logging.debug(f"Published message {self.messages_sent}/{msg_amount}")
    else:
        logging.info(f"Publish complete")

说明

这种方式下,只有当MQTT broker确认了上一条QoS2消息(完成四次握手),才会发送下一条,能保证消息的可靠交付,但间隔会受网络延迟影响,无法严格保证20Hz。


注意事项

  • 不管用哪种方案,都要确保loop_start()启动的后台线程一直运行,不要提前调用loop_stop(),否则QoS2的报文交互会中断,导致消息无法正确确认。
  • 如果你的脚本需要处理其他MQTT消息(比如begin_client和finish_client的指令),后台网络循环会自动处理,不会被定时发布任务阻塞。

备注:内容来源于stack exchange,提问作者Luis Gaspar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.23 10:10:31