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

paho-mqtt回调使用类self属性报Client无write_api属性错误如何解决

MQTT回调AttributeError问题解决

错误根源

你的报错来自三个核心问题:

  1. 回调绑定方式错误:你将on_message绑定到了data_handler类的未绑定方法,而非当前实例的绑定方法。paho-mqtt触发回调时传入的第一个参数是paho自身的Client实例,会被方法接收为self,这个对象自然没有你定义的write_api属性。
  2. 回调参数不匹配:paho-mqtt的on_message回调约定传入3个参数(client, userdata, message),你定义的方法少了userdata参数,会导致参数传递错位。
  3. 监听逻辑错误:反复启停loop_start/loop_stop会导致大量消息丢失,且4秒后就停止监听,不符合长期接收MQTT消息的需求。
    另外你原来直接传入[topic, payload]的写入格式也不符合InfluxDB的写入要求,会触发后续写入报错。

修复要点

  • 回调绑定改为当前实例的方法:self.mqtt_client.on_message = self.mqtt_message
  • 补全回调方法的userdata参数
  • 调整监听逻辑,所有主题订阅完成后仅启动一次事件循环
  • 修正InfluxDB写入数据的格式,使用官方要求的Point结构构造数据

修正后完整代码

# 导入模块
from datetime import datetime
import time
# InfluxDB 依赖
from influxdb_client import InfluxDBClient, Point, WritePrecision
from influxdb_client.client.write_api import SYNCHRONOUS
# MQTT 依赖
import paho.mqtt.client as mqtt

class data_handler():
    def __init__(self, namespace_list=["ESP01","ESP02","ESP03","ESP04","ESP05","ESP06","ESP07","ESP08"]):
        # 初始化InfluxDB客户端
        token = "XXXXXXXXXX"
        self.org = "Home"   
        self.bucket = "HomeSensors"
        self.flux_client = InfluxDBClient(url="http://localhost:8086", token=token)
        self.write_api = self.flux_client.write_api(write_options=SYNCHRONOUS)
        
        # 初始化MQTT客户端
        broker_address="XXX.XXX.XXX.XXX"
        self.mqtt_client = mqtt.Client("influx_client")
        # 绑定当前实例的回调方法
        self.mqtt_client.on_message = self.mqtt_message
        self.mqtt_client.connect(broker_address)

        self.namespace_list = namespace_list
        print(self.namespace_list)

    # 补全userdata参数
    def mqtt_message(self, client, userdata, message):
        payload = str(message.payload.decode("utf-8"))
        print("message received " ,payload)
        print("message topic=",message.topic)
        print("message qos=",message.qos)
        print("message retain flag=",message.retain)
        
        # 构造符合InfluxDB要求的Point数据
        point = Point("sensor_data")\
            .tag("topic", message.topic)\
            .field("value", float(payload))\
            .time(datetime.utcnow(), WritePrecision.NS)
        self.write_api.write(self.bucket, self.org, point)

    def mqtt_listener(self):
        # 先订阅所有主题
        for namespace in self.namespace_list:
            topic = namespace+"/#"
            self.mqtt_client.subscribe(topic, 0)
            print(f"已订阅主题: {topic}")
        # 仅启动一次事件循环,后台持续运行
        self.mqtt_client.loop_start()

def main():
    influxHandler = data_handler(["ESP07"])
    influxHandler.mqtt_listener()
    # 保持主进程不退出,按需调整等待逻辑
    while True:
        time.sleep(1)

if  __name__ == '__main__':
    main()

内容的提问来源于stack exchange,提问作者Martin Jensen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 11:45:01