三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

从零搭建私有物联网平台:Node.js+MQTT+InfluxDB实战指南

从零搭建私有物联网平台:Node.js+MQTT+InfluxDB实战指南

1. 项目概述:为什么我们要从零搭建物联网平台?

如果你正在看这篇文章,大概率和我几年前一样,对“物联网”这个概念既兴奋又迷茫。兴奋的是,万物互联的蓝图听起来就充满未来感;迷茫的是,当你想亲手搭建一个属于自己的物联网平台,用来管理家里的智能设备、监控一个小型温室,或者做一个毕业设计时,面对“平台”、“云端”、“协议”这些词,往往不知从何下手。市面上的阿里云物联网平台、OneNET等公有云服务固然强大,但就像租房子,虽然省心,却难以触及底层,无法完全按你的想法定制规则、处理数据,更关键的是,学习成本被隐藏了,你很难真正理解数据从设备到屏幕的完整旅程。

所以,这个系列教程的目的,就是带你亲手“盖房子”。我们将从最基础的一砖一瓦开始,搭建一个功能完整、架构清晰、完全可控的私有物联网平台。这不是一个简单的Demo,而是一个具备设备接入、数据上报、指令下发、实时监控和简单规则引擎的实战项目。通过这个过程,你将彻底搞懂:一个物联网平台的核心组件有哪些?MQTT、HTTP这些协议到底在通信中扮演什么角色?数据从传感器产生,到最终存入数据库并展示在网页上,中间经历了哪些关键环节?当你理解了这些,再去看那些成熟的公有云平台,就会有一种“哦,原来它是这样工作的”豁然开朗感。

本教程将采用最主流、最易获取的技术栈:后端使用Node.js(得益于其异步高并发特性,非常适合处理海量设备连接),通信协议主打MQTT(物联网领域的事实标准),数据库选用MySQL(存储设备元数据和历史数据)和InfluxDB(专为时序数据优化,存储传感器读数),前端用Vue.js构建一个直观的管理面板。整个项目会部署在一台你能够掌控的云服务器或本地电脑上。放心,每一步我都会详细解释为什么选这个,以及具体怎么操作,哪怕你之前只是写过简单的Python脚本,也能跟着走下来。让我们开始这场从零到一的构建之旅吧。

2. 物联网平台核心架构设计思路

在动手写代码之前,我们必须像建筑师看蓝图一样,先搞清楚我们要建的“房子”长什么样,各个“房间”(模块)如何布局,它们之间如何“走动”(通信)。一个最小化可行物联网平台的核心架构,通常包含以下五个关键部分,我画了一个简单的逻辑图来帮助理解:

[物联网设备] <--(MQTT/HTTP)--> [通信与接入层] <--(内部API)--> [业务逻辑层] <--(读写)--> [数据存储层] | V [Web管理后台] (用户操作界面)

2.1 通信与接入层:设备的“前台接待”

这是平台与物理世界对话的窗口。设备(比如温湿度传感器、智能开关)需要一种方式连接到平台。这里我们主要采用两种协议:

  1. MQTT协议:这是我们的主力。你可以把它想象成“物联网界的微信”。设备(发布者)把数据(消息)发送到一个特定的“话题”(Topic),比如device/001/temperature,平台(订阅者)只要订阅了这个话题,就能立刻收到消息。它的优点是极其轻量、省电、支持海量连接,非常适合传感器间歇性上报数据的场景。我们将使用MosquittoEMQX这类专业的MQTT代理服务器(Broker)来实现。
  2. HTTP/HTTPS协议:作为补充。有些设备能力有限,或者场景需要(比如上传图片),HTTP协议更通用。设备通过POST请求将数据发送到平台指定的API接口。它的优点是简单、易调试,但开销比MQTT大。

这一层的核心职责是协议解析、设备身份认证(鉴权)和连接管理。我们需要判断连接上来的设备是否合法,给它分配一个唯一的内部标识。

2.2 业务逻辑层:平台的“大脑”与“中枢神经”

这是所有魔法发生的地方。它接收来自接入层的原始数据,进行加工、判断和分发。主要包含以下服务:

  • 设备影子服务:这是物联网中一个非常重要的概念。它为每个物理设备在云端创建一个“虚拟镜像”(影子)。设备上报的状态(如“开关:开”)会更新到这个影子里;同时,当你想控制设备时,只需向这个影子发送期望状态(如“开关:关”),业务逻辑层会负责将差异同步给物理设备。这解决了设备离线、网络不稳定时指令无法送达的问题。
  • 规则引擎:这是实现自动化的核心。你可以配置一些“如果...那么...”的规则。例如:“如果设备A的温度>30度,那么向设备B发送打开风扇的指令”。规则引擎会实时处理数据流,触发相应的动作。
  • 数据预处理与转发:将接入层收到的原始数据(可能是一条JSON字符串)进行清洗、格式校验,然后转发给数据存储层进行持久化,或者推送到前端实现实时展示。

这一层我们使用Node.js来构建,利用其事件驱动、非阻塞I/O的特性,高效处理高并发的数据流。

2.3 数据存储层:平台的“记忆仓库”

不同类型的数据,要用不同的“柜子”来存放:

  • 关系型数据库 (MySQL/PostgreSQL):存放结构化、需要关联查询的数据。例如:
    • 设备元信息表:设备ID、名称、型号、所属位置、创建时间。
    • 用户与权限表:哪些用户可以管理哪些设备。
    • 规则定义表:存储上面提到的规则引擎配置。
  • 时序数据库 (InfluxDB/TDengine):专门为时间序列数据优化,这类数据特点是:数据按时间顺序产生、写入多读取少、经常按时间范围聚合查询(如“查询过去24小时的平均温度”)。传感器上报的温度、湿度、电量等数据,天生就是时序数据。用时序数据库来存,读写效率比用MySQL高几个数量级。

2.4 Web管理后台:用户的“控制台与仪表盘”

这是一个给管理员或终端用户使用的网页界面。用Vue.js等前端框架开发,主要功能包括:

  • 设备管理:列表展示、添加/删除设备、查看设备详情。
  • 数据可视化:以图表形式展示设备的历史数据(折线图、仪表盘)。
  • 实时监控:通过WebSocket与后端保持长连接,设备数据一上报,前台页面就实时更新。
  • 规则配置:提供一个可视化界面来创建和管理自动化规则。
  • 指令下发:用户点击网页上的按钮,后台通过API调用,最终将指令下发给设备。

2.5 为什么选择这套技术栈?

  • Node.js:对于IO密集型的物联网应用(大量网络连接、数据库操作),其异步特性优势明显,社区生态繁荣,有丰富的MQTT、数据库客户端库。
  • MQTT:行业标准,几乎所有硬件模块都支持,资源消耗极低。
  • MySQL:最流行的关系型数据库,资料多,易于理解。
  • InfluxDB:时序数据库中的佼佼者,API简单,与Grafana等可视化工具集成完美。
  • Vue.js:渐进式框架,学习曲线平缓,能快速构建出交互丰富的管理界面。

这个架构清晰地将关注点分离,每一层各司其职,便于后续的扩展和维护。比如,未来流量大了,我们可以单独对MQTT Broker集群化,或者对业务逻辑层做微服务拆分。

3. 基础环境搭建与核心组件安装

理论清晰了,现在开始动手准备我们的“施工场地”。我假设你使用一台Ubuntu 20.04 LTS的云服务器或本地虚拟机作为部署环境。这是目前非常稳定且资料丰富的Linux发行版。我们将依次安装所有必需的组件。

3.1 系统更新与基础工具

首先,通过SSH连接到你的服务器,进行系统更新并安装一些后续会用到的工具。

# 更新软件包列表并升级现有软件 sudo apt update && sudo apt upgrade -y # 安装常用工具:vim编辑器、网络工具、压缩解压等 sudo apt install -y vim curl wget net-tools htop unzip git

3.2 安装Node.js与npm

我们使用NodeSource提供的最新LTS版本,比Ubuntu自带的版本要新。

# 下载并执行NodeSource安装脚本(以Node.js 18.x LTS为例) curl -fsSL https://deb.nodesource.com/setup_18.x | sudo -E bash - # 安装Node.js和npm sudo apt install -y nodejs # 验证安装 node --version # 应输出 v18.x.x npm --version # 应输出 9.x.x 或更高

注意:国内服务器访问npm官方源可能较慢,建议立即配置淘宝镜像源,这会极大提升后续安装包的速度。

npm config set registry https://registry.npmmirror.com

3.3 安装与配置MySQL数据库

我们将使用MySQL 8.0来存储设备元数据和用户信息。

# 安装MySQL服务器 sudo apt install -y mysql-server # 启动MySQL服务并设置开机自启 sudo systemctl start mysql sudo systemctl enable mysql # 运行安全安装脚本,进行初始配置 sudo mysql_secure_installation

运行安全脚本时,会提示你:

  1. 设置root密码(务必设置一个强密码并牢记)。
  2. 移除匿名用户?输入Y
  3. 禁止root远程登录?输入N(为了方便后续从本地机器管理,我们允许root远程登录,但会限制IP,生产环境请务必使用更安全的方式)。
  4. 移除测试数据库?输入Y
  5. 重新加载权限表?输入Y

接下来,我们需要进行一些额外配置,允许root用户远程连接(仅限学习测试环境,生产环境强烈建议创建专用用户)。

# 以root身份登录MySQL sudo mysql -u root -p # 输入你刚才设置的密码 -- 在MySQL命令行中执行以下SQL语句 -- 切换到mysql系统数据库 USE mysql; -- 将root用户的认证插件改为mysql_native_password(部分旧客户端兼容需要),并允许从任何主机连接 UPDATE user SET plugin='mysql_native_password', Host='%' WHERE User='root'; -- 刷新权限使更改生效 FLUSH PRIVILEGES; -- 退出 EXIT;

现在,你需要编辑MySQL的配置文件,注释掉绑定本地地址的配置,以允许远程连接。

sudo vim /etc/mysql/mysql.conf.d/mysqld.cnf

找到bind-address这一行(通常在文件中部):

bind-address = 127.0.0.1

将其注释掉或改为0.0.0.0

# bind-address = 127.0.0.1 # 或者 bind-address = 0.0.0.0

保存退出后,重启MySQL服务:

sudo systemctl restart mysql

最后,在你的本地电脑服务器防火墙上,开放MySQL的默认端口3306。

# 如果使用ufw防火墙(Ubuntu默认) sudo ufw allow 3306/tcp sudo ufw reload

3.4 安装与配置InfluxDB

我们将安装InfluxDB 2.x版本,它包含了全新的UI和API。

# 下载并添加InfluxData的APT仓库密钥 wget -q https://repos.influxdata.com/influxdata-archive.key sudo gpg --yes --batch --import influxdata-archive.key # 添加仓库 echo "deb [signed-by=/usr/share/keyrings/influxdata-archive-keyring.gpg] https://repos.influxdata.com/debian stable main" | sudo tee /etc/apt/sources.list.d/influxdata.list # 更新并安装InfluxDB sudo apt update sudo apt install -y influxdb2 # 启动服务并设置开机自启 sudo systemctl start influxdb sudo systemctl enable influxdb

安装完成后,InfluxDB服务会在本地的8086端口启动。我们需要通过命令行进行初始设置。

# 执行设置命令,按照提示操作 influx setup

它会交互式地让你输入:

  • Username:管理员用户名,例如iotadmin
  • Password:设置一个强密码。
  • Organization:组织名称,例如MyIoT
  • Bucket:存储桶名称,相当于数据库,例如sensor_data
  • Retention Period:数据保留期限,输入0表示永久保留(测试用),生产环境请根据需求设置,如52w(52周)。

设置完成后,会生成一个Token,务必妥善保存这个Token,它是后续API访问的凭证。你可以通过以下命令查看配置:

influx auth list

同样,如果需要从外部访问InfluxDB的UI(默认端口8086),记得在防火墙开放该端口。

sudo ufw allow 8086/tcp sudo ufw reload

3.5 安装MQTT代理服务器(Mosquitto)

Mosquitto是一个轻量级、开源的消息代理,完美符合我们的需求。

# 安装Mosquitto Broker和客户端工具 sudo apt install -y mosquitto mosquitto-clients # 启动服务并设置开机自启 sudo systemctl start mosquitto sudo systemctl enable mosquitto

Mosquitto默认使用1883端口(非加密)和8883端口(SSL加密)。为了安全,我们稍后会配置用户名密码认证。现在,可以先测试一下MQTT服务是否正常。

打开两个SSH终端窗口。

在第一个窗口,订阅一个测试主题:

mosquitto_sub -h localhost -t "test/topic" -v

-h指定主机,-t指定主题,-v表示打印详细消息(包括主题名)。

在第二个窗口,向同一个主题发布一条消息:

mosquitto_pub -h localhost -t "test/topic" -m "Hello MQTT!"

如果一切正常,你会在第一个订阅窗口看到输出:test/topic Hello MQTT!。按Ctrl+C可以停止订阅。

至此,我们所有的基础设施——运行环境、数据库、消息队列——都已经准备就绪。下一章,我们将开始编写平台的后端核心逻辑,让这些组件真正联动起来。

4. 后端核心服务开发:构建业务逻辑层

现在,我们进入最核心的编码环节。我们将创建一个Node.js项目,实现设备连接管理、数据接收、存储和规则引擎的雏形。请在你的服务器或开发机上创建一个项目目录。

4.1 项目初始化与依赖安装

# 创建项目目录并进入 mkdir iot-platform-backend && cd iot-platform-backend # 初始化npm项目,一路回车即可 npm init -y # 安装核心依赖 npm install express mqtt mysql2 influxdb-client dotenv cors # 安装开发依赖(用于热重载等) npm install --save-dev nodemon

让我们解释一下这些包的作用:

  • express: Web应用框架,用于提供HTTP API接口。
  • mqtt: MQTT客户端库,让我们的Node.js服务既能订阅设备消息,也能向设备发布消息。
  • mysql2: 高性能的MySQL驱动,用于连接和操作MySQL数据库。
  • influxdb-client: InfluxDB 2.x的官方JavaScript客户端。
  • dotenv: 从.env文件加载环境变量,便于管理配置。
  • cors: 处理跨域资源共享,方便前端调用API。
  • nodemon: 开发工具,监听文件变化自动重启服务。

创建项目入口文件app.js和配置文件.env

touch app.js .env

编辑.env文件,填入你的配置信息:

# 服务器端口 SERVER_PORT=3000 # MySQL 配置 MYSQL_HOST=localhost MYSQL_PORT=3306 MYSQL_USER=root MYSQL_PASSWORD=your_mysql_root_password # 替换为你的密码 MYSQL_DATABASE=iot_platform # MQTT Broker 配置 MQTT_BROKER_URL=mqtt://localhost:1883 # InfluxDB 配置 INFLUXDB_URL=http://localhost:8086 INFLUXDB_TOKEN=your_influxdb_token # 替换为你的Token INFLUXDB_ORG=MyIoT # 替换为你的组织名 INFLUXDB_BUCKET=sensor_data # 替换为你的存储桶名

4.2 数据库表结构设计与初始化

在MySQL中创建我们需要的数据库和表。首先登录MySQL:

mysql -u root -p

执行以下SQL语句:

-- 创建数据库 CREATE DATABASE IF NOT EXISTS iot_platform CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; USE iot_platform; -- 设备表:存储设备基本信息 CREATE TABLE devices ( id INT AUTO_INCREMENT PRIMARY KEY, device_id VARCHAR(64) NOT NULL UNIQUE COMMENT '设备唯一标识,如IMEI、MAC地址', device_name VARCHAR(128) NOT NULL COMMENT '设备名称', device_type VARCHAR(64) COMMENT '设备类型,如temperature_sensor, smart_plug', status ENUM('online', 'offline', 'inactive') DEFAULT 'inactive' COMMENT '设备状态', last_seen TIMESTAMP NULL DEFAULT NULL COMMENT '最后上线时间', created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, INDEX idx_device_id (device_id), INDEX idx_status (status) ) ENGINE=InnoDB COMMENT='设备信息表'; -- 设备影子表:存储设备的期望状态和上报状态 CREATE TABLE device_shadows ( id INT AUTO_INCREMENT PRIMARY KEY, device_id VARCHAR(64) NOT NULL UNIQUE COMMENT '关联设备ID', reported_state JSON COMMENT '设备上报的最新状态,JSON格式', desired_state JSON COMMENT '云端期望的设备状态,JSON格式', metadata JSON COMMENT '元数据,如版本、时间戳', created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, FOREIGN KEY (device_id) REFERENCES devices(device_id) ON DELETE CASCADE, INDEX idx_device_id (device_id) ) ENGINE=InnoDB COMMENT='设备影子表'; -- 规则表:存储自动化规则 CREATE TABLE rules ( id INT AUTO_INCREMENT PRIMARY KEY, rule_name VARCHAR(128) NOT NULL COMMENT '规则名称', condition TEXT NOT NULL COMMENT '触发条件,如 JSONPath 或表达式', action TEXT NOT NULL COMMENT '执行动作,如发布MQTT消息、调用HTTP接口', enabled BOOLEAN DEFAULT TRUE COMMENT '是否启用', created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, INDEX idx_enabled (enabled) ) ENGINE=InnoDB COMMENT='自动化规则表'; -- 插入一条测试设备数据 INSERT INTO devices (device_id, device_name, device_type, status) VALUES ('test_device_001', '客厅温湿度传感器', 'temperature_humidity_sensor', 'inactive'); INSERT INTO device_shadows (device_id, reported_state, desired_state) VALUES ('test_device_001', '{"temperature": 25.5, "humidity": 60}', '{}');

4.3 核心服务模块编写

我们将代码模块化,创建几个核心文件。

4.3.1 数据库连接模块 (db/mysql.js)

const mysql = require('mysql2/promise'); // 使用Promise版本的驱动 require('dotenv').config(); const pool = mysql.createPool({ host: process.env.MYSQL_HOST, port: process.env.MYSQL_PORT, user: process.env.MYSQL_USER, password: process.env.MYSQL_PASSWORD, database: process.env.MYSQL_DATABASE, waitForConnections: true, connectionLimit: 10, // 连接池大小 queueLimit: 0 }); // 提供一个简单的查询函数 async function query(sql, params) { const [rows] = await pool.execute(sql, params); return rows; } module.exports = { pool, query };

4.3.2 InfluxDB客户端模块 (db/influxdb.js)

const { InfluxDB, Point } = require('@influxdata/influxdb-client'); require('dotenv').config(); const influxDB = new InfluxDB({ url: process.env.INFLUXDB_URL, token: process.env.INFLUXDB_TOKEN }); // 获取写入API和查询API const writeApi = influxDB.getWriteApi(process.env.INFLUXDB_ORG, process.env.INFLUXDB_BUCKET); const queryApi = influxDB.getQueryApi(process.env.INFLUXDB_ORG); // 写入时序数据的函数 function writeSensorData(measurement, tags, fields) { const point = new Point(measurement); // 添加标签(索引字段,用于分组查询) for (const [key, value] of Object.entries(tags)) { point.tag(key, value); } // 添加字段(实际存储的数值) for (const [key, value] of Object.entries(fields)) { point.floatField(key, value); } // 写入点 writeApi.writePoint(point); // 确保数据被发送 writeApi.flush().catch(console.error); } module.exports = { writeApi, queryApi, writeSensorData };

4.3.3 MQTT客户端与服务模块 (services/mqttService.js)这是整个后端的数据入口,负责连接MQTT Broker,订阅设备主题,并处理收到的消息。

const mqtt = require('mqtt'); const { query } = require('../db/mysql'); const { writeSensorData } = require('../db/influxdb'); require('dotenv').config(); class MqttService { constructor() { this.client = null; this.connect(); } connect() { // 连接到MQTT Broker this.client = mqtt.connect(process.env.MQTT_BROKER_URL); this.client.on('connect', () => { console.log('✅ MQTT Client Connected to Broker'); // 订阅所有设备上报数据的主题,使用通配符 `+` this.client.subscribe('device/+/data', (err) => { if (!err) { console.log('Subscribed to topic: device/+/data'); } }); // 订阅设备状态主题 this.client.subscribe('device/+/status'); }); this.client.on('message', async (topic, message) => { // 消息是Buffer,需要转成字符串 const msgStr = message.toString(); console.log(`Received message on ${topic}: ${msgStr}`); try { const payload = JSON.parse(msgStr); // 从主题中提取设备ID,例如 topic: "device/abc123/data" -> deviceId: "abc123" const deviceId = topic.split('/')[1]; if (topic.endsWith('/data')) { // 处理传感器数据 await this.handleSensorData(deviceId, payload); } else if (topic.endsWith('/status')) { // 处理设备状态(上线、下线) await this.handleDeviceStatus(deviceId, payload); } } catch (error) { console.error('Error processing MQTT message:', error, 'Raw message:', msgStr); } }); this.client.on('error', (err) => { console.error('MQTT Client Error:', err); }); } // 处理传感器数据:1. 更新设备影子 2. 存入时序数据库 3. 触发规则检查 async handleSensorData(deviceId, data) { console.log(`Handling sensor data from ${deviceId}:`, data); // 1. 更新MySQL中的设备影子(reported_state) const reportedState = JSON.stringify(data); await query( 'UPDATE device_shadows SET reported_state = ?, updated_at = NOW() WHERE device_id = ?', [reportedState, deviceId] ); // 2. 更新设备最后在线时间 await query( 'UPDATE devices SET last_seen = NOW(), status = ? WHERE device_id = ?', ['online', deviceId] ); // 3. 存入InfluxDB时序数据库 // 假设数据格式为 { temperature: 25.5, humidity: 60 } writeSensorData( 'sensor_reading', // 测量名称 { device_id: deviceId }, // 标签 data // 字段 ); // 4. 触发规则引擎(这里先简单打印,后续完善) console.log(`[Rule Engine] Checking rules for device: ${deviceId}`); // TODO: 查询规则表,判断当前数据是否满足任何条件,并执行相应动作 } // 处理设备状态 async handleDeviceStatus(deviceId, status) { const deviceStatus = status.status; // 假设payload为 { "status": "online" } console.log(`Device ${deviceId} status changed to: ${deviceStatus}`); // 更新设备状态 await query( 'UPDATE devices SET status = ?, last_seen = NOW() WHERE device_id = ?', [deviceStatus, deviceId] ); } // 向设备发布消息(用于下发指令) publish(topic, message) { if (this.client && this.client.connected) { const payload = typeof message === 'object' ? JSON.stringify(message) : message; this.client.publish(topic, payload, { qos: 1 }, (err) => { // qos: 1 确保至少送达一次 if (err) { console.error('Publish error:', err); } else { console.log(`Message published to ${topic}: ${payload}`); } }); } else { console.error('MQTT client not connected, cannot publish.'); } } } // 导出单例实例 module.exports = new MqttService();

4.3.4 HTTP API模块 (routes/api.js)提供RESTful API供前端调用,管理设备、查询数据、下发指令。

const express = require('express'); const router = express.Router(); const { query } = require('../db/mysql'); const mqttService = require('../services/mqttService'); // 获取设备列表 router.get('/devices', async (req, res) => { try { const devices = await query('SELECT * FROM devices ORDER BY created_at DESC'); res.json({ code: 0, message: 'success', data: devices }); } catch (error) { console.error(error); res.status(500).json({ code: -1, message: 'Database query failed' }); } }); // 获取单个设备详情及影子 router.get('/device/:deviceId', async (req, res) => { const { deviceId } = req.params; try { const [device] = await query('SELECT * FROM devices WHERE device_id = ?', [deviceId]); if (!device) { return res.status(404).json({ code: -1, message: 'Device not found' }); } const [shadow] = await query('SELECT * FROM device_shadows WHERE device_id = ?', [deviceId]); device.shadow = shadow || {}; res.json({ code: 0, message: 'success', data: device }); } catch (error) { console.error(error); res.status(500).json({ code: -1, message: 'Database query failed' }); } }); // 向设备发送指令(更新设备影子期望状态,并通过MQTT下发) router.post('/device/:deviceId/command', async (req, res) => { const { deviceId } = req.params; const { command } = req.body; // 例如 { "power": "on" } if (!command || typeof command !== 'object') { return res.status(400).json({ code: -1, message: 'Invalid command format' }); } try { // 1. 更新影子中的期望状态 const desiredState = JSON.stringify(command); await query( 'INSERT INTO device_shadows (device_id, desired_state) VALUES (?, ?) ON DUPLICATE KEY UPDATE desired_state = ?', [deviceId, desiredState, desiredState] ); // 2. 通过MQTT向设备发布指令 // 主题格式:device/{deviceId}/command const topic = `device/${deviceId}/command`; mqttService.publish(topic, command); res.json({ code: 0, message: 'Command sent successfully' }); } catch (error) { console.error(error); res.status(500).json({ code: -1, message: 'Failed to send command' }); } }); // 查询设备历史数据(从InfluxDB查询,这里先预留接口,具体查询在下一部分实现) router.get('/device/:deviceId/history', async (req, res) => { res.json({ code: 0, message: 'History data endpoint (to be implemented with InfluxDB query)' }); }); module.exports = router;

4.3.5 主应用文件 (app.js)将上述模块整合起来。

const express = require('express'); const cors = require('cors'); require('dotenv').config(); // 初始化服务(MQTT连接会自动建立) require('./services/mqttService'); const apiRoutes = require('./routes/api'); const app = express(); const PORT = process.env.SERVER_PORT || 3000; // 中间件 app.use(cors()); // 允许跨域 app.use(express.json()); // 解析JSON请求体 app.use(express.urlencoded({ extended: true })); // 路由 app.use('/api', apiRoutes); // 健康检查端点 app.get('/health', (req, res) => { res.json({ status: 'UP', service: 'IoT Platform Backend' }); }); // 启动服务器 app.listen(PORT, () => { console.log(`🚀 IoT Platform Backend server is running on http://localhost:${PORT}`); });

现在,一个具备设备数据接收、存储、影子管理和基础API的物联网平台后端核心就搭建完成了。你可以使用nodemon来启动开发服务器:

npx nodemon app.js

如果看到🚀 IoT Platform Backend server is running...✅ MQTT Client Connected to Broker的日志,说明服务启动成功。

5. 平台功能测试与数据流验证

后端跑起来了,但我们还需要验证整个数据链路是否通畅。我们将模拟一个物联网设备,从数据上报到存储,再到前端查询的完整过程。

5.1 模拟设备数据上报

我们使用mosquitto_pub命令行工具来模拟设备。假设设备ID为test_device_001

1. 模拟设备上线:

mosquitto_pub -h localhost -t "device/test_device_001/status" -m '{"status": "online"}'

查看后端服务日志,应该能看到Device test_device_001 status changed to: online的提示。同时,查询MySQL数据库,devices表中对应设备的status应变为onlinelast_seen更新为当前时间。

2. 模拟上报温湿度数据:

mosquitto_pub -h localhost -t "device/test_device_001/data" -m '{"temperature": 26.8, "humidity": 55}'

后端日志应显示处理了传感器数据。此时检查:

  • MySQL:执行SELECT * FROM device_shadows WHERE device_id='test_device_001';reported_state字段应更新为{"temperature": 26.8, "humidity": 55}
  • InfluxDB:我们可以通过命令行查询刚写入的数据。
    # 使用InfluxDB CLI查询 influx query ' from(bucket:"sensor_data") |> range(start: -1h) |> filter(fn: (r) => r._measurement == "sensor_reading" and r.device_id == "test_device_001") |> yield(name: "last_hour") '
    你应该能看到两条记录,分别对应temperaturehumidity字段及其值、时间戳。

5.2 测试HTTP API

使用curl或 Postman 测试我们编写的API。

1. 获取设备列表:

curl -X GET http://localhost:3000/api/devices

应返回包含test_device_001的设备列表JSON。

2. 获取特定设备详情:

curl -X GET http://localhost:3000/api/device/test_device_001

应返回该设备的详细信息及其影子状态。

3. 向设备发送指令:

curl -X POST http://localhost:3000/api/device/test_device_001/command \ -H "Content-Type: application/json" \ -d '{"command": {"led": "on"}}'

后端日志应显示Message published to device/test_device_001/command: {"led":"on"}。同时,检查MySQL中device_shadows表,desired_state字段应更新为{"led":"on"}

5.3 模拟设备接收指令

在另一个终端,模拟设备订阅自己的指令主题:

mosquitto_sub -h localhost -t "device/test_device_001/command" -v

保持这个订阅窗口打开。然后再次执行上面的curl命令发送指令。你会在订阅窗口立刻看到收到的指令消息:device/test_device_001/command {"led":"on"}

至此,我们验证了物联网平台最核心的数据流:

  1. 设备 -> 平台 (上行):设备通过MQTT发布状态和数据到Broker,后端服务订阅并处理,分别存入MySQL(元数据、影子)和InfluxDB(时序数据)。
  2. 平台 -> 设备 (下行):用户通过HTTP API发起指令,后端更新影子并通过MQTT发布到特定主题,设备订阅该主题并接收指令。

这个闭环的打通,意味着我们的物联网平台“心脏”已经开始跳动。在下一篇教程中,我们将构建一个现代化的Vue.js前端管理界面,实现设备可视化、实时数据图表和规则配置,让这个平台真正“有血有肉”。我们还会完善规则引擎,并探讨如何增强安全性(如MQTT TLS加密、JWT API认证)和部署优化。

← 返回列表