MQTT转InfluxDB仅存入0值问题求助
问题排查与修复方案
核心问题分析
1. MQTT主题完全不匹配(最关键)
发布端的温度主题是tp/eng/iotp_project/temp、湿度主题是tp/eng/iotp_project/hum,但订阅端订阅的是tp/iotp_project/temperature——两者路径不一致,导致订阅端根本接收不到任何数据,temp变量一直保持初始值0.0,最终写入InfluxDB的全是0。
2. 订阅端变量作用域错误
订阅端main函数内的dbclient = InfluxDBClient(...)是局部变量,而外部while循环调用的是全局的dbclient = None,实际执行write_points时会触发异常,但因无捕获逻辑,程序要么崩溃要么静默失败,你看到的0值是程序启动时写入的初始值。
3. 发布端传感器数据未实时更新
发布端的temp和hum仅在程序启动时采集一次,后续循环调用MQTT()时始终发送固定初始值,不会同步最新温湿度。
4. 订阅端消息处理冗余延迟
onMessage回调里的time.sleep(2)完全多余,会阻塞MQTT消息线程,导致数据接收不及时。
修复后的完整代码
发布端代码
import RPi.GPIO as GPIO import time import smbus2 import bme280 import paho.mqtt.client as mqtt mqtt_broker = "test.mosquitto.org" topic_temp = "tp/eng/iotp_project/temp" topic_hum = "tp/eng/iotp_project/hum" port = 1 address = 0x76 bus = smbus2.SMBus(port) calibration_params = bme280.load_calibration_params(bus, address) # GPIO配置 fan = 23 buzzer = 22 fan_threshold = 24 buzzer_threshold = 80 GPIO.setmode(GPIO.BCM) GPIO.setup(fan, GPIO.OUT) GPIO.setup(buzzer, GPIO.OUT) # 初始化MQTT长连接 my_mqtt = mqtt.Client() my_mqtt.connect(mqtt_broker, port=1883) print("--Connected to MQTT broker") def read_and_control_sensor(): # 实时采集传感器数据 sensor = bme280.sample(bus, address, calibration_params) temp = round(sensor.temperature, 2) hum = round(sensor.humidity, 2) # 风扇/蜂鸣器逻辑 GPIO.output(fan, GPIO.HIGH if temp > fan_threshold else GPIO.LOW) GPIO.output(buzzer, GPIO.HIGH if hum > buzzer_threshold else GPIO.LOW) return temp, hum try: while True: temp, hum = read_and_control_sensor() # 发布最新数据 my_mqtt.publish(topic_temp, temp) my_mqtt.publish(topic_hum, hum) print(f"--Published: Temp={temp}C, Hum={hum}%") time.sleep(2) finally: # 程序退出时清理资源 my_mqtt.disconnect() GPIO.cleanup() GPIO.setwarnings(False)
订阅端/InfluxDB代码
from influxdb import InfluxDBClient import time import paho.mqtt.client as mqtt PASSWORD = 'root' USER = 'root' DBNAME = 'temperature' HOST = 'localhost' mqtt_broker = "test.mosquitto.org" topic_temp = "tp/eng/iotp_project/temp" topic_hum = "tp/eng/iotp_project/hum" PORT = 8086 dbclient = None current_temp = 0.0 current_hum = 0.0 def onMessage(client, userdata, message): global current_temp, current_hum payload = float(message.payload.decode()) # 根据主题区分温湿度 if message.topic == topic_temp: current_temp = payload print(f"--Received Temp: {current_temp}C") elif message.topic == topic_hum: current_hum = payload print(f"--Received Hum: {current_hum}%") def startMQTT(): my_mqtt = mqtt.Client() my_mqtt.on_message = onMessage my_mqtt.connect(mqtt_broker, port=1883) my_mqtt.subscribe(topic_temp) my_mqtt.subscribe(topic_hum) my_mqtt.loop_start() def get_db_data_point(): now = time.gmtime() return [ { "time": time.strftime("%Y-%m-%d %H:%M:%S", now), "measurement": 'reading', "tags": {"nodeId": "node_1"}, "fields": {"Temperature": current_temp, "Humidity": current_hum} } ] def main(): global dbclient startMQTT() # 初始化数据库连接并验证 try: dbclient = InfluxDBClient(HOST, PORT, USER, PASSWORD, DBNAME) dbclient.ping() print("--Connected to InfluxDB") except Exception as e: print(f"--InfluxDB connection failed: {e}") return while True: data_point = get_db_data_point() try: dbclient.write_points(data_point) print(f"--Written to DB: Temp={current_temp}, Hum={current_hum}") except Exception as e: print(f"--DB write error: {e}") time.sleep(2) if __name__ == '__main__': main()
关键修改说明
- 主题统一:订阅端使用和发布端完全一致的主题,确保数据能正常接收。
- 变量作用域修正:订阅端将
dbclient声明为全局变量,避免局部变量导致的空指针错误。 - 实时数据采集:发布端新增
read_and_control_sensor()函数,每次循环重新采集传感器数据,保证发布的是最新值。 - MQTT长连接优化:发布端保持MQTT长连接,避免频繁创建/断开连接的开销。
- 异常捕获:添加数据库连接和写入的异常捕获,方便快速定位问题。
内容的提问来源于stack exchange,提问作者BorisLimTP
相关产品推荐
相关产品推荐

