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

无法通过订阅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脚本无响应的核心原因有两个:

  1. MQTT客户端未正确配置:缺少on_connect、on_message回调函数,也没有启动客户端循环,导致脚本运行后没有持续监听MQTT消息,直接执行完毕跳转到下一个单元格。
  2. 实时绘图逻辑不完整:没有初始化数据存储结构,也没有在收到消息后更新图表并刷新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("正在停止绘图...")
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 12:55:28