paho-mqtt回调使用类self属性报Client无write_api属性错误如何解决
MQTT回调AttributeError问题解决
错误根源
你的报错来自三个核心问题:
- 回调绑定方式错误:你将
on_message绑定到了data_handler类的未绑定方法,而非当前实例的绑定方法。paho-mqtt触发回调时传入的第一个参数是paho自身的Client实例,会被方法接收为self,这个对象自然没有你定义的write_api属性。 - 回调参数不匹配:paho-mqtt的
on_message回调约定传入3个参数(client, userdata, message),你定义的方法少了userdata参数,会导致参数传递错位。 - 监听逻辑错误:反复启停
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
相关产品推荐
相关产品推荐

