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

调用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()

代码说明

  1. 用latest_msg变量保存最新收到的MQTT消息,旧消息会被直接覆盖,确保只处理最新的图像。
  2. processing标志位避免同时处理多条消息,防止重复执行延迟逻辑。
  3. 使用loop_start()启动客户端线程,配合主循环定时检查并处理消息,避免阻塞MQTT客户端的消息接收。

内容的提问来源于stack exchange,提问作者HAL 9000

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 18:15:50