调用reinitialise()后触发Invalid Host错误的paho MQTT订阅端问题
问题描述
尝试调用reinitialise()方法重置缓冲区,以此在降低帧率的同时保留实时图像(用time.sleep(0.4)模拟后处理延迟),但调用该方法后出现ValueError: Invalid host错误。
订阅端代码
# set broker MQTT_SERVER = "localhost" #Write Server IP Address MQTT_PATH = "test" client = paho.Client() client.connect(MQTT_SERVER, 1883, 60) def on_connect(client, userdata, flags, rc): print("Connected with result code "+str(rc)) client.subscribe(MQTT_PATH) def on_message(client, userdata, msg): imageStream = io.BytesIO(msg.payload) imageFile = Image.open(imageStream) cv2.imshow('Client',np.array(imageFile)) cv2.waitKey(1) time.sleep(0.4) client.reinitialise() client.on_connect = on_connect client.on_message = on_message client.loop_forever()
报错堆栈信息
ValueError Traceback (most recent call last) <ipython-input-4-02d8506b1627> in <module> 14 client.on_connect = on_connect 15 client.on_message = on_message ---> 16 client.loop_forever() c:\users\wikto\appdata\local\programs\python\python38\lib\site-packages\paho\mqtt\client.py in loop_forever(self, timeout, max_packets, retry_first_connection) 1777 else: 1778 try: -> 1779 self.reconnect() 1780 except (OSError, WebsocketConnectionError): 1781 self._handle_on_connect_fail() c:\users\wikto\appdata\local\programs\python\python38\lib\site-packages\paho\mqtt\client.py in reconnect(self) 1014 connect()/connect_async().""" 1015 if len(self._host) == 0: -> 1016 raise ValueError('Invalid host.') 1017 if self._port <= 0: 1018 raise ValueError('Invalid port number.') ValueError: Invalid host.
问题原因
reinitialise()方法会重置MQTT客户端的所有内部状态,包括存储的主机地址、端口等连接信息。当客户端在loop_forever()中自动尝试重连时,发现_host字段为空,就会抛出Invalid host错误。且该方法完全不适合用来实现“降低帧率并保留实时图像”的需求。
解决方案
要实现降低帧率同时只显示最新图像,应该通过丢弃旧消息、只处理最新消息的方式,而非重置客户端。以下是修改后的代码:
import paho.mqtt.client as paho import io from PIL import Image import cv2 import numpy as np import time MQTT_SERVER = "localhost" MQTT_PATH = "test" client = paho.Client() client.connect(MQTT_SERVER, 1883, 60) # 标记是否正在处理消息,以及保存最新收到的消息 processing = False latest_msg = None def on_connect(client, userdata, flags, rc): print(f"Connected with result code {rc}") client.subscribe(MQTT_PATH) def process_image(): global processing, latest_msg if processing or latest_msg is None: return processing = True # 处理最新消息 image_stream = io.BytesIO(latest_msg.payload) image_file = Image.open(image_stream) cv2.imshow('Client', np.array(image_file)) cv2.waitKey(1) # 模拟后处理延迟 time.sleep(0.4) # 处理完成后清空消息,允许接收新的最新消息 latest_msg = None processing = False def on_message(client, userdata, msg): global latest_msg # 直接覆盖旧消息,只保留最新的一条 latest_msg = msg client.on_connect = on_connect client.on_message = on_message # 启动客户端线程 client.loop_start() try: while True: # 定时检查并处理最新消息 process_image() time.sleep(0.01) except KeyboardInterrupt: # 清理资源 client.loop_stop() cv2.destroyAllWindows()
代码说明
- 用
latest_msg变量保存最新收到的MQTT消息,旧消息会被直接覆盖,确保只处理最新的图像。 processing标志位避免同时处理多条消息,防止重复执行延迟逻辑。- 使用
loop_start()启动客户端线程,配合主循环定时检查并处理消息,避免阻塞MQTT客户端的消息接收。
内容的提问来源于stack exchange,提问作者HAL 9000
相关产品推荐
相关产品推荐

