无法通过订阅MQTT Broker绘制BME688传感器数据实时图表
无法通过订阅MQTT Broker绘制BME688传感器数据实时图表
问题背景
我正在基于Bosch BME688气味传感器阵列和Adafruit ESP32 Huzzah Feather开发板做项目,已经完成了传感器采集数据并推送至Mosquitto MQTT Broker的功能——在Mac终端订阅MQTT主题能看到实时数据,Arduino串口监视器也能正常输出完整数据。但在Jupyter Notebook中编写的订阅MQTT并实时绘图的脚本完全无响应,运行后直接跳转到下一个单元格,没有任何图表输出,明明传感器和Broker都在正常工作。
正常工作的Arduino MQTT发布代码
#include <bsec2.h> #include "commMux.h" #include "mqtt_datalogger.h" #include <WiFi.h> #include <PubSubClient.h> #define NUM_OF_SENS 8 #define PANIC_LED LED_BUILTIN #define ERROR_DUR 1000 const char* ssid = "mywifiname"; const char* password = "mywifipassword"; const char* mqttServer = "MyMQTTServerAddress"; const int mqttPort = "MySensorsPort"; const char* mqttTopic = "sensorData"; const char* mqttClientName = "MyClientName"; WiFiClient espClient; PubSubClient mqttClient(espClient); // create MQTT logger bme68xData sensorData[NUM_OF_SENS] = {0}; mqttDataLogger logger(&mqttClient, NUM_OF_SENS, mqttTopic); void reconnect() { while (!mqttClient.connected()) { Serial.print("Attempting MQTT connection..."); if (mqttClient.connect(mqttClientName)) { Serial.println("connected"); } else { Serial.print("failed, rc="); Serial.print(mqttClient.state()); Serial.println(" try again in 5 seconds"); delay(5000); } } } void errLeds(void); void checkBsecStatus(Bsec2 bsec); void newDataCallback(const bme68xData data, const bsecOutputs outputs, Bsec2 bsec); Bsec2 envSensor[NUM_OF_SENS]; comm_mux communicationSetup[NUM_OF_SENS]; uint8_t bsecMemBlock[NUM_OF_SENS][BSEC_INSTANCE_SIZE]; uint8_t sensor = 0; void setup() { bsecSensor sensorList[] = { BSEC_OUTPUT_IAQ, BSEC_OUTPUT_RAW_TEMPERATURE, BSEC_OUTPUT_RAW_PRESSURE, BSEC_OUTPUT_RAW_HUMIDITY, BSEC_OUTPUT_RAW_GAS, BSEC_OUTPUT_STABILIZATION_STATUS, BSEC_OUTPUT_RUN_IN_STATUS }; Serial.begin(115200); comm_mux_begin(Wire, SPI); pinMode(PANIC_LED, OUTPUT); delay(100); while(!Serial) delay(10); Serial.print("Connecting to "); Serial.print(ssid); Serial.println("..."); WiFi.begin(ssid, password); while (WiFi.status() != WL_CONNECTED) { delay(500); Serial.print("."); } Serial.println(""); Serial.println("WiFi connected"); Serial.println("IP address: "); Serial.println(WiFi.localIP()); mqttClient.setServer(mqttServer, mqttPort); mqttClient.setBufferSize(600); Serial.print("MQTT client buffer size: "); Serial.println(mqttClient.getBufferSize()); reconnect(); logger.beginSensorData(); for (uint8_t i = 0; i < NUM_OF_SENS; i++) { communicationSetup[i] = comm_mux_set_config(Wire, SPI, i, communicationSetup[i]); envSensor[i].allocateMemory(bsecMemBlock[i]); if (!envSensor[i].begin(BME68X_SPI_INTF, comm_mux_read, comm_mux_write, comm_mux_delay, &communicationSetup[i])) { checkBsecStatus (envSensor[i]); } if (!envSensor[i].updateSubscription(sensorList, ARRAY_LEN(sensorList), BSEC_SAMPLE_RATE_LP)) { checkBsecStatus (envSensor[i]); } envSensor[i].attachCallback(newDataCallback); } Serial.println("BSEC library version " + \ String(envSensor[0].version.major) + "." \ + String(envSensor[0].version.minor) + "." \ + String(envSensor[0].version.major_bugfix) + "." \ + String(envSensor[0].version.minor_bugfix)); } void loop() { for (sensor = 0; sensor < NUM_OF_SENS; sensor++) { if (!envSensor[sensor].run()) { checkBsecStatus(envSensor[sensor]); } } reconnect(); } void errLeds(void) { while(1) { digitalWrite(PANIC_LED, HIGH); delay(ERROR_DUR); digitalWrite(PANIC_LED, LOW); delay(ERROR_DUR); } } void newDataCallback(const bme68xData data, const bsecOutputs outputs, Bsec2 bsec) { if (!outputs.nOutputs) { return; } Serial.println("BSEC outputs:\n\tsensor num = " + String(sensor)); Serial.println("\ttimestamp = " + String((int) (outputs.output[0].time_stamp / INT64_C(1000000)))); for (uint8_t i = 0; i < outputs.nOutputs; i++) { const bsecData output = outputs.output[i]; switch (output.sensor_id) { case BSEC_OUTPUT_IAQ: Serial.println("\tiaq = " + String(output.signal)); Serial.println("\tiaq accuracy = " + String((int) output.accuracy)); break; case BSEC_OUTPUT_RAW_TEMPERATURE: Serial.println("\ttemperature = " + String(output.signal)); sensorData[sensor].temperature = output.signal; break; case BSEC_OUTPUT_RAW_PRESSURE: Serial.println("\tpressure = " + String(output.signal)); sensorData[sensor].pressure = output.signal; break; case BSEC_OUTPUT_RAW_HUMIDITY: Serial.println("\thumidity = " + String(output.signal)); sensorData[sensor].humidity = output.signal; break; case BSEC_OUTPUT_RAW_GAS: Serial.println("\tgas resistance = " + String(output.signal)); sensorData[sensor].gas_resistance = output.signal; break; case BSEC_OUTPUT_STABILIZATION_STATUS: Serial.println("\tstabilization status = " + String(output.signal)); break; case BSEC_OUTPUT_RUN_IN_STATUS: Serial.println("\trun in status = " + String(output.signal)); break; default: break; } } logger.assembleAndPublishSensorData(sensor, &sensorData[sensor]); } void checkBsecStatus(Bsec2 bsec) { if (bsec.status < BSEC_OK) { Serial.println("BSEC error code : " + String(bsec.status)); errLeds(); } else if (bsec.status > BSEC_OK) { Serial.println("BSEC warning code : " + String(bsec.status)); } if (bsec.sensor.status < BME68X_OK) { Serial.println("BME68X error code : " + String(bsec.sensor.status)); errLeds(); } else if (bsec.sensor.status > BME68X_OK) { Serial.println("BME68X warning code : " + String(bsec.sensor.status)); } }
收到的MQTT数据示例
{ "datapoints" : [ 00:00:52.829 -> [ 0, 37272, 21.483234, 9.934442, 33.501423, 37847.425781 ], 00:00:52.829 -> [ 1, 37328, 21.446659, 9.932629, 33.433811, 16072.325195 ], 00:00:52.829 -> [ 2, 37383, 21.332052, 9.936784, 33.085041, 20506.248047 ], 00:00:52.861 -> [ 3, 37438, 21.383812, 9.933786, 33.250164, 27663.712891 ], 00:00:52.861 -> [ 4, 37493, 21.828791, 9.933929, 32.243370, 26661.111328 ], 00:00:52.861 -> [ 5, 37548, 21.818996, 9.935399, 31.534727, 22519.353516 ], 00:00:52.861 -> [ 6, 37603, 22.092115, 9.936502, 32.235497, 23791.822266 ], 00:00:52.861 -> [ 7, 37658, 22.005573, 9.934257, 32.204479, 35605.007812 ] 00:00:52.861 -> ] }
注:数据点格式为[传感器ID, 时间戳, 温度, 压力, 湿度, 气体电阻]
问题原因分析
你的Jupyter脚本无响应的核心原因有两个:
- MQTT客户端未正确配置:缺少
on_connect、on_message回调函数,也没有启动客户端循环,导致脚本运行后没有持续监听MQTT消息,直接执行完毕跳转到下一个单元格。 - 实时绘图逻辑不完整:没有初始化数据存储结构,也没有在收到消息后更新图表并刷新Matplotlib的交互界面。
修正后的Jupyter实时绘图代码
import paho.mqtt.client as mqtt import matplotlib.pyplot as plt import numpy as np import time import json # 优先设置Notebook交互后端(比plt.ion()更稳定) %matplotlib widget # MQTT配置:注意端口号必须是整数类型 mqtt_server = "MyMQTTServerAddress" mqtt_port = 1883 # 替换为你的实际MQTT端口,默认是1883 mqtt_topic = "sensorData" # 初始化数据存储:为8个传感器分别记录各数据项的历史 data = {} for sensor_id in range(8): data[sensor_id] = { 'timestamps': [], 'temperature': [], 'pressure': [], 'humidity': [], 'gas_resistance': [] } # 初始化Matplotlib多子图实时绘图 fig, axes = plt.subplots(2, 2, figsize=(14, 10)) axes = axes.flatten() sensor_colors = plt.cm.viridis(np.linspace(0, 1, 8)) # 为8个传感器分配不同颜色 plot_lines = {} # 配置温度子图 axes[0].set_title('Temperature (°C)') axes[0].set_xlabel('Timestamp') axes[0].set_ylabel('Temperature') for sensor_id in range(8): line, = axes[0].plot([], [], color=sensor_colors[sensor_id], label=f'Sensor {sensor_id}') plot_lines[f'temp_{sensor_id}'] = line axes[0].legend(fontsize=8) # 配置湿度子图 axes[1].set_title('Humidity (%)') axes[1].set_xlabel('Timestamp') axes[1].set_ylabel('Humidity') for sensor_id in range(8): line, = axes[1].plot([], [], color=sensor_colors[sensor_id]) plot_lines[f'hum_{sensor_id}'] = line # 配置压力子图 axes[2].set_title('Pressure (hPa)') axes[2].set_xlabel('Timestamp') axes[2].set_ylabel('Pressure') for sensor_id in range(8): line, = axes[2].plot([], [], color=sensor_colors[sensor_id]) plot_lines[f'press_{sensor_id}'] = line # 配置气体电阻子图 axes[3].set_title('Gas Resistance (Ω)') axes[3].set_xlabel('Timestamp') axes[3].set_ylabel('Gas Resistance') for sensor_id in range(8): line, = axes[3].plot([], [], color=sensor_colors[sensor_id]) plot_lines[f'gas_{sensor_id}'] = line # MQTT连接成功回调:订阅目标主题 def on_connect(client, userdata, flags, rc): if rc == 0: print("已连接到MQTT Broker,开始订阅主题...") client.subscribe(mqtt_topic) else: print(f"连接失败,错误码: {rc}") # MQTT消息接收回调:处理数据并更新图表 def on_message(client, userdata, msg): try: # 处理带串口时间戳前缀的payload:截取JSON部分 payload_str = msg.payload.decode('utf-8') if "->" in payload_str: payload_str = payload_str.split("->")[-1].strip() # 解析JSON数据 payload_json = json.loads(payload_str) datapoints = payload_json.get('datapoints', []) # 逐个处理传感器数据点 for point in datapoints: if len(point) != 6: print(f"数据格式错误,跳过: {point}") continue sensor_id, timestamp, temp, press, hum, gas = point sensor_id = int(sensor_id) # 存储新数据 data[sensor_id]['timestamps'].append(timestamp) data[sensor_id]['temperature'].append(temp) data[sensor_id]['pressure'].append(press) data[sensor_id]['humidity'].append(hum) data[sensor_id]['gas_resistance'].append(gas) # 更新对应传感器的绘图线 plot_lines[f'temp_{sensor_id}'].set_data(data[sensor_id]['timestamps'], data[sensor_id]['temperature']) plot_lines[f'hum_{sensor_id}'].set_data(data[sensor_id]['timestamps'], data[sensor_id]['humidity']) plot_lines[f'press_{sensor_id}'].set_data(data[sensor_id]['timestamps'], data[sensor_id]['pressure']) plot_lines[f'gas_{sensor_id}'].set_data(data[sensor_id]['timestamps'], data[sensor_id]['gas_resistance']) # 自动调整坐标轴范围,避免数据溢出显示区域 for ax in axes: ax.relim() ax.autoscale_view() # 刷新图表 fig.canvas.draw() fig.canvas.flush_events() time.sleep(0.1) # 控制刷新频率,避免界面卡顿 except Exception as e: print(f"处理消息出错: {e},原始payload: {msg.payload.decode('utf-8')}") # 创建并配置MQTT客户端 client = mqtt.Client() client.on_connect = on_connect client.on_message = on_message # 连接MQTT Broker并启动消息循环 try: client.connect(mqtt_server, mqtt_port, 60) client.loop_start() # 后台线程处理MQTT消息 print("客户端已启动,等待接收数据...") except Exception as e: print(f"无法连接到MQTT Broker: {e}") # 保持脚本运行,按Ctrl+C停止 try: while True: time.sleep(1) except KeyboardInterrupt: print("正在停止绘图...")
相关产品推荐
相关产品推荐

