物联网项目最让人头疼的从来不是"连不上",而是设备一多、数据一杂,整个链路就开始互相拖后腿。我做过好几个从传感器到看板的完整项目,最典型的一个场景是:两百多个采集节点,每秒上报一次温湿度和电流数据,前端还要实时刷新曲线。一开始图省事,MQTT 消息直接往 MongoDB 里塞,结果跑了三天,查询慢到打不开页面,磁盘也快被写满。后来把架构拆成"消息接入层 + 双存储 + 可视化层",才算真正稳住。这篇就把这套从零搭建的物联网数据中台完整拆一遍,涉及 MQTT 协议接入、Node.js 服务编排、MongoDB 与 InfluxDB 双存储分工,以及 React 前端可视化。不管你是刚接触物联网后端,还是已经写过几个 Demo 想往生产级靠,这套结构都能直接拿去改。
1. 先想清楚数据中台到底在解决什么问题
很多人一上来就问"用哪个 MQTT 服务器""MongoDB 怎么装",其实顺序反了。数据中台的核心不是某个组件,而是把"设备产生数据"到"人看到数据"这条链路拆成职责清晰的几段,每段用最合适的工具。想不清楚这一点,后面选型全是拍脑袋。
1.1 物联网数据的三个天然特征
设备数据和我们平时做的业务数据完全不是一回事,它有三个绕不开的特征,直接决定了架构长什么样。
第一是写多读少且写入极其频繁。一个中等规模的采集场景,几百个设备每秒上报,一天就是几千万条记录。这种量级下,如果用传统关系型数据库,光是写入就会把连接池打满。第二是数据带时间戳且天然按时间查询。你几乎不会去查"某条具体的温度记录",而是查"过去一小时的平均温度""昨天下午三点的电流峰值",时间永远是查询的第一维度。第三是冷热数据价值差异巨大。最近几小时的数据要秒级可查、要能画实时曲线,但三个月前的原始数据基本只有归档和偶尔回溯的价值。
这三个特征合在一起,就解释了为什么单一存储方案一定会翻车。MongoDB 擅长存结构灵活、需要按业务字段查询的数据,比如设备元信息、告警记录、配置快照;InfluxDB 这类时序数据库天生为"按时间写入、按时间聚合"优化,压缩率高、聚合查询快。把它们混用,才是这套架构的关键。
1.2 双存储分工的边界在哪里
我见过不少人纠结"到底用 MongoDB 还是 InfluxDB",其实这个问题本身就问错了。正确的问法是:哪类数据放哪个库。下面这张表是我在实际项目里总结的分工原则,可以直接对照。
| 数据类型 | 存储选择 | 原因 |
|---|---|---|
| 设备元信息(型号、位置、固件版本) | MongoDB | 字段不固定,需要按业务条件查询 |
| 实时遥测数据(温度、电流、湿度) | InfluxDB | 高频写入,按时间聚合查询 |
| 告警事件与处理记录 | MongoDB | 需要关联业务状态、支持复杂条件筛选 |
| 长期趋势与统计指标 | InfluxDB | 降采样后存储,压缩比高 |
| 用户配置与权限 | MongoDB | 结构化、低频读写 |
注意:不要把原始遥测数据同时写进两个库"做备份",这会让写入放大一倍,收益却几乎为零。真正需要的是各司其职,而不是冗余。
1.3 整体链路的分层设计
把链路拆开看,从下到上大致是四层。设备接入层负责 MQTT 连接、主题订阅、消息解析;服务编排层用 Node.js 做消息消费、数据清洗、路由分发;存储层是 MongoDB 和 InfluxDB 双写;可视化层用 React 拉取聚合数据并渲染图表。每一层之间通过明确的接口通信,任何一层出问题都不会直接拖垮其他层。
这种分层最大的好处是可替换。比如哪天 MQTT 服务器要换,只要接入层的接口不变,上层完全无感;前端要换框架,后端 API 也不用动。做中台最怕的就是"牵一发动全身",分层就是给未来留退路。
2. MQTT 接入层:主题设计和消息解析的坑
MQTT 是整个链路的入口,它轻量、基于发布订阅、适合弱网环境,这些优点大家都知道。但真正决定项目能不能扩展的,是主题(Topic)怎么设计、消息格式怎么定。这两件事一旦定错,后期改起来要动所有设备,代价极大。
2.1 主题层级设计要一次到位
MQTT 的主题是用斜杠分隔的层级结构,支持通配符订阅。设计主题时,我建议遵循"从粗到细、从稳定到易变"的原则。一个经过验证的格式是:
{产品线}/{设备类型}/{设备ID}/{数据类型}举个例子,factory-a/sensor/dev-001/temperature。这样设计的好处是,服务端可以用factory-a/sensor/+/temperature订阅所有传感器的温度数据,也可以用factory-a/#订阅整个厂区的所有消息。通配符+匹配单层,#匹配多层,这是 MQTT 协议里非常实用的能力。
提示:设备 ID 一定要放在靠后的层级,因为它是变化最频繁的部分。如果把易变字段放在前面,通配符订阅会变得非常别扭。
2.2 消息体格式:JSON 是默认答案,但别乱塞
消息体我强烈建议用 JSON,可读性好、各语言解析都方便。但很多人会把整个设备状态一股脑塞进一条消息,导致单条消息几 KB,高频上报时带宽和解析开销都上去了。我的做法是一条消息只承载一个语义单元,比如温度就是温度,电流就是电流,需要批量上报时用数组包一层。
{ "deviceId": "dev-001", "ts": 1718000000000, "value": 23.6, "unit": "celsius" }时间戳字段ts一定要带,而且用毫秒级 Unix 时间戳。不要依赖服务端接收时间,因为网络延迟和批量上报会让接收时间和真实采集时间差出好几秒,画曲线时就会错位。
2.3 Node.js 侧的消息消费与解析
Node.js 处理 MQTT 消息,最常用的是mqtt这个库。它的异步模型非常适合 IO 密集的消息消费场景。核心逻辑其实不复杂:连接服务器、订阅主题、在message事件里解析并分发。
const mqtt = require('mqtt'); const client = mqtt.connect('mqtt://broker-host:1883', { clientId: `ingest-${process.pid}`, clean: false, reconnectPeriod: 2000 }); client.on('connect', () => { client.subscribe('factory-a/+/+/+', { qos: 1 }, (err) => { if (err) console.error('订阅失败', err); }); }); client.on('message', async (topic, payload) => { try { const data = JSON.parse(payload.toString()); await routeMessage(topic, data); } catch (e) { console.error('消息解析失败', topic, e.message); } });这里有几个细节值得说。clientId带上进程号,是为了多实例部署时不会互相顶掉连接。clean: false配合 QoS 1,能在断线重连后尽量不丢消息。reconnectPeriod设成 2 秒,避免网络抖动时疯狂重连。
2.4 QoS 等级怎么选才不浪费
MQTT 有三个 QoS 等级:0 是"最多一次",1 是"至少一次",2 是"恰好一次"。很多人一上来就用 QoS 2,觉得最安全,其实代价很大——QoS 2 需要四次握手,吞吐量会明显下降。
我的经验是:普通遥测数据用 QoS 0 或 1,关键指令用 QoS 1,几乎不用 QoS 2。遥测数据丢一两条对趋势分析影响很小,用 QoS 0 换吞吐量更划算;告警类消息不能丢,用 QoS 1;QoS 2 的"恰好一次"在物联网场景里收益有限,反而拖慢整体速度。
3. Node.js 服务编排:把消息变成可查询的数据
接入层拿到消息只是第一步,真正决定数据质量的是服务编排层。这一层要做三件事:数据清洗、路由分发、双写存储。听起来简单,但每一件都有讲究。
3.1 数据清洗:脏数据比没数据更可怕
设备上报的数据经常不干净:温度偶尔冒出 999、时间戳是 0、字段缺失。如果这些脏数据直接进库,后面画出来的曲线会突然跳一下,聚合统计也会被污染。所以清洗这一步不能省。
我通常做三类校验:范围校验(温度是否在 -50 到 150 之间)、时间戳校验(是否在合理时间窗口内,比如不能是未来时间)、必填字段校验(deviceId、ts、value 是否齐全)。任何一项不过,直接丢弃并记一条日志,不要试图"修复"。
function validate(data) { if (!data.deviceId || typeof data.value !== 'number') return false; if (data.ts <= 0 || data.ts > Date.now() + 60000) return false; if (data.unit === 'celsius' && (data.value < -50 || data.value > 150)) return false; return true; }注意:校验规则要可配置,不同设备类型的合理范围不一样。硬编码在代码里,加一种新设备就得改代码重新部署,很麻烦。
3.2 路由分发:按数据类型决定去向
清洗通过后,数据要根据类型决定写到哪里。遥测数据进 InfluxDB,设备状态变更、告警进 MongoDB。这个路由逻辑最好抽成一个独立的模块,用配置驱动,而不是写一堆 if-else。
const routes = [ { match: (t) => t.endsWith('/temperature') || t.endsWith('/current'), target: 'influx' }, { match: (t) => t.endsWith('/status') || t.endsWith('/alarm'), target: 'mongo' } ]; function resolveTarget(topic) { const route = routes.find((r) => r.match(topic)); return route ? route.target : null; }这样加新数据类型时,只要加一条路由规则,不用动核心逻辑。我在一个项目里就是靠这个设计,从最初的两种数据类型扩展到十几种,核心代码一行没改。
3.3 批量写入:性能的关键在这里
单条写入是性能杀手。InfluxDB 和 MongoDB 都支持批量写入,把短时间内的数据攒一批再写,吞吐量能提升一个数量级。我的做法是用一个内存队列,攒够 500 条或者超过 1 秒就触发一次批量写。
const buffer = []; let timer = null; function enqueue(point) { buffer.push(point); if (buffer.length >= 500) { flush(); } else if (!timer) { timer = setTimeout(flush, 1000); } } async function flush() { if (timer) { clearTimeout(timer); timer = null; } if (buffer.length === 0) return; const batch = buffer.splice(0, buffer.length); await influx.writePoints(batch); }这个"攒批"逻辑看着简单,但要注意两点:一是进程退出前必须把缓冲区刷干净,否则会丢最后一批数据;二是批量写失败要有重试,不能直接丢掉。
3.4 背压处理:别让内存被消息撑爆
高频场景下,如果写入速度跟不上消息到达速度,内存队列会越堆越大,最后 OOM。这就是背压问题。解决办法是给队列设上限,超过上限时要么丢弃最旧的数据,要么暂停消费。
我一般用"丢弃最旧 + 记录丢弃计数"的策略,因为对遥测数据来说,最新的数据永远比旧数据有价值。同时把丢弃计数暴露成一个监控指标,一旦持续增长就说明写入能力不足,该扩容了。
4. 双存储落地:MongoDB 与 InfluxDB 各管一摊
存储层是这套架构的核心。MongoDB 和 InfluxDB 的定位完全不同,用对了事半功倍,用错了就是给自己挖坑。这一章把两者的安装、建模、查询都过一遍。
4.1 MongoDB 的安装与建模要点
MongoDB 的安装本身不难,但新手经常卡在几个地方。Windows 上装的时候,最容易忽略的是数据目录和日志目录要手动创建,否则服务起不来。另外安装时如果勾选了"作为服务运行",记得配置文件的路径要对。
建模方面,MongoDB 是文档型数据库,不需要预先定义表结构,但这不代表可以随便存。我的原则是:同一集合内的文档结构尽量一致,这样索引才有效。设备元信息可以这样存:
{ _id: ObjectId("..."), deviceId: "dev-001", type: "temperature-sensor", location: { factory: "factory-a", line: "line-1" }, firmware: "1.2.3", installedAt: ISODate("2024-01-15T00:00:00Z"), status: "online" }索引是 MongoDB 性能的命脉。deviceId这种高频查询字段一定要建索引,否则数据量一大,查询就会全表扫描。我见过一个项目因为没建索引,几百万条数据查一次要好几秒。
db.devices.createIndex({ deviceId: 1 }, { unique: true }); db.alarms.createIndex({ deviceId: 1, createdAt: -1 });4.2 InfluxDB 的数据模型与写入
InfluxDB 的数据模型和传统数据库差别很大,它由 measurement(类似表)、tag(带索引的标签)、field(实际数值)、timestamp 组成。理解这个模型是高效使用的前提。
关键原则是:tag 用来存需要按它筛选和分组的维度,field 用来存真正的数值。比如设备 ID、设备类型适合做 tag,温度值适合做 field。因为 tag 会被索引,查询快,但 tag 的基数(不同值的数量)不能太高,否则索引会膨胀。
const { InfluxDB, Point } = require('@influxdata/influxdb-client'); const point = new Point('telemetry') .tag('deviceId', 'dev-001') .tag('type', 'temperature') .floatField('value', 23.6) .timestamp(Date.now() * 1e6); // 纳秒提示:InfluxDB 的时间戳默认是纳秒精度,从 JavaScript 的毫秒时间戳转换时要乘以 1e6,这个细节不注意会导致数据时间错乱。
4.3 查询对比:什么时候用哪个库
两个库的查询能力差异很大,选错了会很痛苦。下面这张表是我实际用下来的对比。
| 查询需求 | 推荐库 | 说明 |
|---|---|---|
| 查某设备最近 1 小时温度曲线 | InfluxDB | 时间范围聚合,原生支持 |
| 查某厂区所有离线设备 | MongoDB | 按业务字段筛选 |
| 计算过去 24 小时平均电流 | InfluxDB | 内置聚合函数 |
| 查某设备的历史告警记录 | MongoDB | 关联业务状态 |
| 按小时降采样长期趋势 | InfluxDB | 连续查询或任务 |
InfluxDB 的聚合查询写起来很直观,比如查过去一小时每分钟的平均温度:
SELECT MEAN("value") FROM "telemetry" WHERE "deviceId" = 'dev-001' AND time > now() - 1h GROUP BY time(1m)MongoDB 则更适合这种带业务条件的查询:
db.alarms.find({ status: 'unresolved', createdAt: { $gte: new Date(Date.now() - 86400000) } }).sort({ createdAt: -1 });4.4 数据保留策略:别让磁盘被撑爆
时序数据会无限增长,必须设置保留策略。InfluxDB 支持在 bucket 级别设置保留期,比如原始数据保留 30 天,降采样后的数据保留 1 年。这样既控制了存储成本,又保留了长期趋势。
MongoDB 这边,告警记录这类数据也要定期归档。我的做法是写一个定时任务,把超过半年的告警记录导出到冷存储,然后从主库删除。别小看这一步,一个跑了半年的项目,告警表很容易涨到几千万条。
5. React 可视化:让数据真正被看见
数据存得再好,看不到等于没有。可视化层是用户唯一直接接触的部分,它的流畅度和准确性直接决定项目口碑。React 生态里有大量图表库,选型和数据处理都有讲究。
5.1 图表库选型:别只看颜值
React 图表库很多,ECharts、Recharts、Chart.js、Visx 各有特点。我的选型逻辑是看三个维度:数据量、交互复杂度、定制需求。
| 图表库 | 适合场景 | 数据量承受力 |
|---|---|---|
| Recharts | 常规业务图表,快速开发 | 中等 |
| ECharts | 复杂交互、大数据量 | 高 |
| Chart.js | 轻量简单图表 | 中等 |
| Visx | 高度定制 | 取决于实现 |
如果是实时曲线、需要缩放拖拽,ECharts 更稳;如果是常规的仪表盘,Recharts 开发效率更高。我一般实时监控页用 ECharts,配置页用 Recharts。
5.2 实时数据刷新:轮询还是推送
前端拿实时数据有两种方式:定时轮询和 WebSocket 推送。轮询实现简单,但延迟高、请求多;WebSocket 实时性好,但要处理断线重连。
我的建议是:秒级刷新用轮询,亚秒级用 WebSocket。大多数工业监控场景,3 到 5 秒刷新一次完全够用,轮询就够了。只有对实时性要求极高的场景才上 WebSocket。
useEffect(() => { const fetchData = async () => { const res = await fetch(`/api/telemetry?deviceId=${id}&range=1h`); setData(await res.json()); }; fetchData(); const timer = setInterval(fetchData, 5000); return () => clearInterval(timer); }, [id]);注意:轮询一定要在组件卸载时清理定时器,否则页面切换后请求还在跑,内存泄漏就是这么来的。
5.3 大数据量渲染的性能陷阱
实时曲线最容易踩的坑是数据点太多导致卡顿。一个设备一小时每分钟一个点才 60 个,但如果原始数据是每秒一个点,一小时就是 3600 个点,多个设备叠加就上万了。浏览器渲染上万个点会明显掉帧。
解决办法是在服务端做降采样,前端只拿聚合后的数据。比如画一小时曲线,服务端返回每分钟的平均值,60 个点足够画出趋势。InfluxDB 的GROUP BY time(1m)就是干这个的。前端不要试图自己处理原始数据,那是自找麻烦。
5.4 状态管理:别让数据流变成一团乱麻
React 项目做大了,状态管理很容易失控。我的经验是:服务端数据用 React Query 这类库管理,本地 UI 状态用 useState/useReducer。不要把服务端返回的数据塞进全局 store,那样缓存、刷新、失效都要自己处理,很容易出 bug。
React Query 自带缓存、自动重试、后台刷新,处理轮询数据特别合适。配置好refetchInterval,它自己就会定时拉数据,组件卸载自动停止,比手写 setInterval 省心得多。
6. 联调与排错:那些文档不会告诉你的坑
架构搭起来只是开始,真正花时间的是联调。这一章把我踩过的坑集中列一下,都是文档里不会写、但实际一定会遇到的。
6.1 时间戳精度不一致导致曲线错位
最常见的问题:设备上报毫秒时间戳,InfluxDB 存纳秒,前端展示时又按秒处理,三层精度不统一,曲线就会错位。解决办法是全链路统一用毫秒时间戳,只在写入 InfluxDB 的那一刻转成纳秒,读取时再转回来。
6.2 时区问题让数据"穿越"
InfluxDB 默认用 UTC 存储,前端展示如果不做时区转换,用户看到的曲线会比实际时间差 8 小时。这个坑我踩过不止一次。正确做法是:存储统一 UTC,展示时按用户时区转换,转换逻辑集中在一个工具函数里,不要散落各处。
6.3 连接池耗尽与重连风暴
Node.js 服务在高并发下,如果 MongoDB 连接池配置太小,会出现请求排队;如果 MQTT 重连策略太激进,断网恢复瞬间会有一堆实例同时重连,把服务器打挂。我的配置是:MongoDB 连接池按并发量设,一般 50 到 100;MQTT 重连加随机抖动,避免同时重连。
reconnectPeriod: 2000 + Math.random() * 10006.4 内存泄漏的排查思路
Node.js 服务跑几天就 OOM,八成是内存泄漏。常见原因有三个:事件监听器没移除、定时器没清理、缓冲区无限增长。排查时可以用process.memoryUsage()定期打印内存,观察是否持续增长。定位到具体模块后,重点检查监听器和定时器的生命周期。
提示:MQTT 客户端的
message事件如果每次都on一次而不off,监听器会越积越多,这是最隐蔽的泄漏点之一。
7. 从 Demo 到生产还差哪几步
能跑通不等于能上线。这套架构从 Demo 到生产,还有几件事必须补上,否则迟早出事。
7.1 监控与告警不能省
生产环境必须监控几个核心指标:消息消费速率、写入延迟、队列积压量、各服务存活状态。任何一个异常都要能及时告警。我一般用 Prometheus 采集指标,Grafana 做看板,简单直接。
7.2 数据一致性校验
双存储最大的风险是两边数据对不上。要定期做一致性校验,比如对比 InfluxDB 里某设备某小时的点数和 MongoDB 里记录的应上报次数。发现偏差要及时排查,可能是消息丢失,也可能是写入失败。
7.3 灰度与回滚能力
任何改动都要能灰度发布、快速回滚。MQTT 主题变更、存储结构变更这类影响面大的操作,一定要先在部分设备上验证,确认没问题再全量。别一次性全推,出事就是全站故障。
7.4 容量规划要提前做
按当前设备数和上报频率,估算每天的写入量和存储增长,提前规划磁盘和扩容方案。我见过太多项目是磁盘写满了才发现,那时候只能紧急清理,手忙脚乱。
这套架构我在几个项目里反复打磨过,最大的体会是:物联网中台的难点从来不在单个技术,而在它们之间的配合。MQTT 的主题设计决定了扩展性,Node.js 的批处理和背压决定了稳定性,双存储的分工决定了查询性能,React 的数据处理决定了用户体验。每一环都要想清楚"为什么这么做",而不是照抄别人的方案。真正跑起来之后你会发现,那些文档里一笔带过的细节,才是决定项目成败的地方。