如何基于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()
关键修改点说明
- 用
loop_start()替代loop_forever():loop_start()会启动一个后台线程专门处理MQTT的网络交互(比如心跳、QoS2的PUBREC/PUBREL报文),主线程可以自由处理定时任务,不会被阻塞。 - 用
threading.Timer实现定时发布:每次发布一条消息后,调度50ms后的下一次发布,避免了阻塞式的sleep循环。 - 用事件保持主线程运行:添加
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
相关产品推荐
相关产品推荐

