十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

自研IoT平台实战:从MQTT接入到设备影子与时序数据架构

自研IoT平台实战:从MQTT接入到设备影子与时序数据架构 最近在技术社区里总有人问一个很扎心的问题IoT Platform 这种东西开源方案一抓一大把AWS IoT Core 这类云服务也很成熟为什么还有人要自己动手从头搭一套说实话我以前也是这个想法直到自己做的智能硬件项目越来越多设备从几十台涨到上千台才发现通用平台的很多合理设计在私有场景里非常别扭。有的平台把每条消息按条计费设备一多成本直接失控有的平台规则引擎太重一个简单的开关控制要写一堆流程。折腾一圈下来我决定自己动手实现一套真正贴合业务的 IoT Platform这也是这个系列 Part 1 的由来。这篇文章会先讲清楚自研平台的必要性和架构思路然后落到底层最核心的三块地基MQTT 设备接入网关、设备影子的状态管理、时序数据链路。读完你能得到一套可以照着落地的从零搭建方案也能看清每一层设计背后的实际考虑。适合正在开发智能硬件、想摆脱第三方平台束缚、或者单纯想了解物联网平台内部原理的开发者。1. 先想清楚自研 IoT 平台的边界在哪里自研平台最忌讳一上来就画大饼。我见过不少团队刚开始雄心勃勃想把 AWS IoT Core 的所有功能都复刻一遍结果大半年过去连设备接入都没跑通。自己的平台核心诉求只有一个字够用。所以在动手之前必须把自己做什么和不做什么这两件事彻底想明白。1.1 什么时候真的需要自研什么时候是重复造轮子我自己的判断标准很简单如果业务只用到设备连接、数据上行、指令下发这三板斧那真的没必要自研市面上的开源平台、云托管服务完全覆盖得了。但当出现下面几种情况自研的价值就开始显现了。第一种情况是数据敏感。智能家居、工业采集、医疗设备这类场景数据出域本身就是大忌讳。设备上报的温度、湿度、电压曲线甚至设备位置放在第三方平台上面即使对方承诺不碰数据你也很难睡得踏实。自研平台把数据牢牢掌握在自己手里合规审查也更容易过关。第二种情况是成本失控。云平台大多按连接时长、消息条数、数据存储量三层叠加计费。设备少的时候一个月几十块钱看不出来等设备规模到万台、消息量到亿级费用曲线会陡到让你怀疑人生。自研之后主要成本变成几台云服务器和带宽费用基本可控。第三种情况是协议定制需求。通用平台遵循的是平台规范你得按它的物模型、规则引擎、告警模板来做。但实际项目里很多设备是私有协议、私有数据格式有的甚至要支持远程升级、远程调试这种平台很少做的功能。到了这一步定制化就成了刚需通用平台反而变成了束缚。1.2 这个系列要解决的问题范围既然定位是自己的平台第一步就要把范围划定清楚。我的规划是分几步走Part 1 先解决设备怎么接进来、状态怎么管、数据怎么存这三个最底层的问题Part 2 处理规则引擎、告警通知、更丰富的数据分析Part 3 再做多租户权限、OTA 固件升级、可视化大屏这些上层应用。Part 1 的范围说得很清楚不碰前端可视化不碰复杂规则链不碰大数据分析。这三个东西虽然也重要但不是平台的地基过早投入会让整个项目陷入泥潭。把设备接入、状态管理、数据存储这三件事做到极致后面的所有上层功能都会轻松很多。1.3 技术选型的核心原则自研平台很容易陷入什么都想用最新最酷的误区。我的原则是越老越稳的用老方案越需要交互的越用简单方案。核心组件选型如下层次方案选型理由设备接入EMQX 开源版MQTT 协议实现成熟单机百万连接无压力支持集群扩展状态管理Redis 6.x读取写入都在微秒级天然适合设备状态的读多写多场景设备台账PostgreSQL 15设备信息、用户绑定关系关系型结构清晰事务能力可靠历史数据InfluxDB 2.x时序写入吞吐极高自带降采样和保留策略省心后端服务Go Gin轻量高并发协程模型天然适合 IoT 场景部署又简单这套组合从成本结构看非常友好全部组件都是开源免费或者社区版够用的水平。唯一需要额外注意的就是 EMQX 的集群 License 问题不过单机承载几万个设备完全没有压力等真正到了需要集群的阶段项目本身带来的收益早就可以覆盖商业授权费用了。2. 平台整体架构一条消息从传感器到数据库的完整旅程IoT 平台看起来模块很多其实核心只有一条数据管道设备上报数据 → 平台接收 → 状态更新 → 历史存储 → 上层消费。设备指令下行则是反向管道。把这条管道画清楚整个平台的架构就清晰了。2.1 核心模块与数据流向整个平台分成四个核心模块分别承担不同的职责接入网关层负责处理设备连接、鉴权、消息收发。这里直接选用 EMQX 作为 MQTT Broker它负责最底层的协议解析和连接管理我们不需要自己写网络层只需要处理业务逻辑。消息处理层消费 MQTT 消息做格式校验、数据清洗然后分别写入 Redis实时状态和 InfluxDB时序数据。同时负责设备指令下发把用户操作变成 MQTT 消息推给设备。业务服务层对外提供 REST API管理设备台账、用户体系、设备生命周期。手机 App 和 Web 管理后端都走这一层。存储层MySQL 存设备档案、PostgreSQL 存业务数据、Redis 存实时状态、InfluxDB 存历史时序。各司其职互不干扰。数据流向大致是这样的设备端通过 MQTT 协议连接接入网关上报的消息经过消息处理层后分两条路走——一条路更新设备影子Redis另一条路写入时序库InfluxDB。当用户需要下发指令时业务服务层通过 API 更新设备影子的期望状态再由指令服务把命令通过 MQTT 推给设备。设备执行完操作后再上报实际状态整个闭环就完成了。2.2 为什么用 MQTT 而不是 HTTP/WebSocketIoT 平台的消息协议业界主流就是 MQTT没有之一。做这个架构选址的时候我认真对比过 HTTP、WebSocket、MQTT 三者在 IoT 场景下的表现结果 MQTT 在关键指标上是碾压级的。先说 HTTP。每台设备如果每隔几秒上报一次数据一分钟 60 个请求一万台设备就是每分钟 60 万请求这个量对服务端压力和带宽消耗都很大。HTTP 是请求-响应模型即使设备一条数据都没得报也得定时发起请求非常浪费。MQTT 则是发布-订阅模型设备与 Broker 之间保持一条长连接数据上报直接推给服务器不需要像 HTTP 那样重复握手建连。尤其大量设备处于休眠状态、偶尔才上报一次数据的场景MQTT 的省电优势非常明显。虽然 WebSocket 也是长连接但它没有 QoS 机制没有遗嘱消息没有订阅管理这些对设备通信来说都是非常头疼的事情。MQTT 的 QoS 0/1/2 三级机制分别对应尽力而为至少一次恰好一次很多 IoT 场景要求数据至少一次不丢失这是 WebSocket 做不到的。2.3 EMQX 作为接入网关的配置思路EMQX 本身是一个很成熟的 MQTT Broker在架构里它是整个平台的咽喉部位配置上要特别注意三个地方。监听端口要区分设备端和内部服务端。设备端走 1883MQTT和 8883MQTT over TLS内部服务端之间走独立的端口还需要配置认证插件。EMQX 提供 HTTP 认证插件可以对接业务服务的鉴权接口实现设备凭证的动态校验。ACL 权限控制同样重要只允许设备订阅自己所属的主题防止设备越权订阅其他设备的数据流。这一点上EMQX 提供了基于客户端的 ACL 规则配置可以利用 clientid 和主题的通配符规则组合实现。最后是集群配置虽然初始阶段单机够用但架构上还是要预留水平扩展能力。EMQX 节点之间通过组播协议自动发现后续需要扩容时只需要再把 EMQX 的集群方式改成基于节点列表的方式不用改任何应用代码就能横向扩展。3. 设备接入层如何设计一套能扛住生产环境的消息网关接入层是整个平台的最前线不夸张地说设备接入设计的好坏直接决定后续平台好不好用。这一节讲的是我在实操中验证过的完整接入方案包含主题规范、数据格式、鉴权方式、心跳检测、离线感知五个关键环节。3.1 主题命名规范前期没做这件事后面会花十倍的代价改MQTT 的主题是一个分层结构设计得好就是一套自文档化的消息路由系统设计得差就成了消息黑洞。我第一次做主题设计的时候用的是类似/device/001/temp这种混乱的写法后来设备类型一多规则引擎匹配起来简直像噩梦。经过几轮实践我沉淀了一套比较稳妥的主题规范分为下行指令、上行数据、设备事件三种类型下行指令iot/{productKey}/{deviceName}/cmd/set 请求方下发 iot/{productKey}/{deviceName}/cmd/reply 设备应答 上行数据iot/{productKey}/{deviceName}/data/report 设备主动上报 iot/{productKey}/{deviceName}/data/reply 平台应答用于确认 设备事件iot/{productKey}/{deviceName}/event/online 上线通知 iot/{productKey}/{deviceName}/event/offline 离线通知主题分层的核心目的是让路由规则可以用通配符批量匹配。例如规则引擎如果想处理所有设备的温度上报可以直接订阅iot///data/report省去逐个设备匹配。再比如业务需要按产品线隔离第一层可以直接用产品 key 来区分规则引擎和权限控制都会轻松很多。3.2 设备数据格式统一物模型是数据质量的命根子设备发的数据如果格式五花八门平台永远只会在接手的边缘疯狂兼容。有的设备上报{temperature:25.6}有的上报{temp:25.6}还有的上报{data:{temperature:25.6}}。规则引擎为了兼容这些会变成一坨浆糊。我的做法是定义一套固定的物模型格式不同设备类型通过物模型模板来约束字段。消息统一用 JSON核心结构如下{ version: 1.0, deviceId: dev-001, timestamp: 1704792012000, properties: { temperature: 25.6, humidity: 60.2 }, events: [ { type: threshold_exceeded, description: 温度超过阈值 } ] }这套格式背后的逻辑是version解决设备端升级时数据版本不同的问题deviceId解决主题里有设备号但消息体里没人验证的安全隐患timestamp使用毫秒时间戳避免设备端和服务器时区不一致导致的混乱properties是核心业务数据events用于承载设备主动上报的事件。规则引擎和存储层都只认这套规范省掉一大部分数据清洗的功夫。3.3 鉴权方式给设备发身份证而不是密码设备接入平台的鉴权最容易犯的错误是想用请求式的交互完成。物联网设备和 App 用户不一样App 用户可以在网页上输账号密码但设备固件里没法弹一个登录框。我采用的是三元组鉴权方案设备接入时携带产品密钥ProductKey、设备名DeviceName、设备密钥DeviceSecret。平台用产品密钥找到该产品的 HMAC 密钥然后使用 HMAC 算法对设备名加当前时间戳做签名与设备传上来的签名进行比较。密钥本身不直接传输只传签名结果可以很好地防止中间人截获。// 设备端签名生成逻辑伪代码 func GenerateSignature(deviceSecret string, deviceName string, timestamp string) string { raw : fmt.Sprintf(%s%s, deviceName, timestamp) h : hmac.New(sha256.New, []byte(deviceSecret)) h.Write([]byte(raw)) return hex.EncodeToString(h.Sum(nil)) }设备初次接入平台时通过产线预置或者云端注册的方式分配三元组即可。平台收到连接请求后调用认证服务完成签名校验通过认证才允许连接。这个鉴权机制不依赖 TLS 时安全性也够用配合 MQTT over TLS 的双层防护整体安全性可以应对大多数场景。3.4 心跳与离线检测让平台知道设备活着还是死了设备断连是 IoT 场景的常态网络抖动、断电、模块重启随时可能发生。平台必须快速感知设备离线才能及时做出告警、推送给用户。MQTT 协议自带心跳机制客户端在 CONNECT 报文里声明 KeepAlive 时间单位秒Broker 如果超过 1.5 倍 KeepAlive 时间没有收到客户端消息就判定连接断开。设备端的心跳间隔我建议设置在 30 到 120 秒之间太短会增加功耗和流量太长会导致平台离线检测滞后。配合心跳机制还需要用到 MQTT 的遗嘱消息Last Will and Testament。设备上线时在 EMQX 里登记一条遗嘱消息如果意外掉线Broker 会自动帮助设备发布这条遗嘱消息。我们在遗嘱消息里携带设备离线事件这样平台可以立刻感知设备掉线而不是等待心跳超时。这里有一个实际的坑要提醒设备主动休眠时有些 SDK 会把 TCP 连接断开但不会清理 MQTT 会话。结果就是平台以为设备还在线设备其实已经睡死过去了。所以设备休眠前一定要显式发送 DISCONNECT 报文并标记 CleanSession 为 true让平台能及时清理会话状态。4. 设备影子让平台永远知道设备现在长什么样IoT 平台里有一个概念叫设备影子Device Shadow它本质上是平台侧维护的一份设备当前状态缓存。所有对设备状态的查询都直接查缓存而不是实时去问设备。为什么要这么做因为分布式系统里实时状态查询的代价太大了。4.1 什么是设备影子为什么需要它没有设备影子的平台想知道某台设备当前温度通常有两种做法要么直接向设备发查询指令等设备返回数据要么依赖设备定时上报从最近一条上报里取数。第一种做法延迟不可控设备断网时根本查不到第二种做法也不靠谱设备几小时不上报这个最近状态早就过期了。设备影子的思路很简单平台侧维护一份状态镜像设备每次上报的数据都实时更新这份镜像查询状态时直接读镜像。设备离线了也能查到最后一次上报的状态这在智能家居控制场景里特别关键。4.2 设备影子的数据结构设计设备影子的数据结构和 Linux 的 proc 文件系统有点像至少包含两部分信息实际状态reported和期望状态desired。{ deviceId: dev-001, version: 42, timestamp: 1704792012000, reported: { temperature: 25.6, humidity: 60.2, switch: on }, desired: { switch: off } }reported是设备上报的实际运行状态平台被动接收。desired是用户或应用下发的期望状态平台主动设置。比如用户在 App 上点了一下关灯应用的指令先写入desired然后指令服务通过 MQTT 把命令推给设备设备执行完后再上报实际状态把reported更新成off同时把desired清除。这套机制的价值在于指令下发与指令执行是解耦的。设备离线时指令可以安全地停留在desired里设备上线后主动来拉取期望状态再执行这就是 IoT 里经典的离线指令能力。4.3 用 Redis 实现设备影子版本号与原子更新实现设备影子的方案很多我选 Redis 是因为它的读写性能极高原子操作也适合状态更新场景。每个设备影子的 Key 采用device:shadow:{deviceId}的命名方式Value 直接存 JSON 字符串。更新影子时需要注意并发问题。设备上报和 App 下发可能同时发生如果直接读取-修改-写回很容易因为并发覆盖导致状态不一致。Redis 的 WATCH/MULTI/EXEC 事务可以解决这个问题更简单的方法是使用版本号。# 更新设备影子的原子逻辑伪代码 import redis import json r redis.Redis(hostlocalhost, port6379) def update_shadow(device_id, patch): key fdevice:shadow:{device_id} with r.pipeline() as pipe: while True: try: pipe.watch(key) shadow json.loads(pipe.get(key) or {}) shadow.update(patch) shadow[version] shadow.get(version, 0) 1 pipe.multi() pipe.set(key, json.dumps(shadow)) pipe.execute() break except redis.WatchError: continue版本号除了保证并发的原子性之外还有一个作用让设备端能感知自己的状态是否被平台纠正。例如设备上报温度 25.6 度但业务侧通过影子把期望状态改成了 26 度。设备端在收到指令时可以看到 version 变化据此判断是否需要重新同步状态。4.4 指令下发的完整闭环讲完设备影子指令下发的完整流程就水到渠成了。用户从 App 点了一下打开空调到空调真实开启中间经历了四个关键环节应用层调用业务 API把期望状态写入设备影子的desired字段。指令服务监听影子变化发现desired有更新立即把指令通过 MQTT 推送给目标设备。设备收到指令后执行动作比如继电器闭合、压缩机启动完成后上报新状态。平台把设备上报的状态更新到影子的reported字段同时清除desired整个闭环完成。这个闭环设计的精妙之处在于第一步和第二步之间出现任何失败都不会产生消息丢失的问题。设备离线的时候指令暂时留在desired里设备上线时平台会把desired下发下去等设备上报后再说。这套机制和云厂商 IoT 平台的设计思路是同一个模型实际跑起来可靠性很高。5. 数据链路从设备上报到可控查询的完整打通设备接入进来了状态也管住了还差最后一块拼图历史数据怎么存、怎么查。设备上报的数据如果没有落到可查询的存储里这个平台就只是联调玩具。这一节我们打通从 MQTT 消息到 InfluxDB 的完整链路并给出可落地的实现方案。5.1 为什么选择 InfluxDB 这样的时序数据库IoT 数据有一个显著特征大部分是时序数据。温度、湿度、电压、功率这些数据是随着时间不断产生的每条数据自带时间戳。传统关系型数据库不是不能存而是存了之后查询、聚合、清理都很难受。举个例子一万台设备每 10 秒上报一条数据一天的记录数就是 8640 万条。在 MySQL 里对这个量级做SELECT AVG(temperature) WHERE device_idxxx AND time 2024-01-01这种查询没有建立合适的索引的话慢查询会把整个库拖垮。InfluxDB 这类时序数据库天生就是为这种模式优化的写入吞吐高、压缩率高、按时间分片查询迅速。5.2 时序数据模型的建模要点InfluxDB 的数据模型是 measurement、tag、field、timestamp 四个维度。IoT 场景的设计规范是measurement 表示设备类型或数据类型tag 存设备的筛选维度设备 ID、产品 key、地理位置field 存实际的数值指标timestamp 就是数据产生的时刻。temperature,deviceIddev-001,productKeyprod-01 value25.6 1704792012000000000 temperature,deviceIddev-002,productKeyprod-01 value24.8 1704792012000000000 humidity,deviceIddev-001,productKeyprod-01 value60.2 1704792012000000000每一行的格式都是measurement 后跟逗号分隔的 tag 列表然后空格再跟 field 列表最后是纳秒时间戳。tag 和 field 的区别要记牢tag 会被索引field 不会被索引所以筛选条件尽量放 tag查询结果才快。5.3 消息处理服务从 MQTT 到 InfluxDB 的桥梁MQTT 消息要写入 InfluxDB中间需要一个消费服务。EMQX 支持通过规则引擎把 MQTT 消息直接转发到 Webhook 或者 Kafka也可以使用桥接方式直接写到 InfluxDB。我选的是 Webhook 转发方式EMQX 规则引擎把消息转发到自建的 Go 服务Go 服务做数据清洗后再批量写入 InfluxDB。这种方式的好处是数据清洗逻辑掌握在自己手里EMQX 只负责转发。比如有的设备上报的湿度值是 0 到 100 的整数有的上报的是 0 到 1 的浮点数清洗服务负责换算成统一单位。规则引擎里还顺便把设备断连事件遗嘱消息转发出来方便做设备离线统计。// 批量写入 InfluxDB 的核心逻辑简化版 func WriteToInfluxDB(messages []DeviceMessage) { points : make([]*influxdb2.Point, 0, len(messages)) for _, msg : range messages { p : influxdb2.NewPoint( telemetry, map[string]string{ deviceId: msg.DeviceID, productKey: msg.ProductKey, }, map[string]interface{}{ temperature: msg.Temperature, humidity: msg.Humidity, }, time.UnixMilli(msg.Timestamp), ) points append(points, p) } writeAPI.WritePoint(context.Background(), points...) }批量写入是时序数据写入的基本功。单条消息逐条写入 InfluxDB 的效率很低但批量打包写入后吞吐量可以轻松达到几万 point/s这个量级对 IoT 平台完全够用。写入时还可以开启 gzip 压缩网络带宽占用会明显下降。5.4 存储策略与降采样数据不能只存不管时序数据是越存越多越存越贵如果不管存储生命周期数据量会变成一笔巨额账单。InfluxDB 提供了保留策略Retention Policy和连续查询Continuous Query配合起来可以实现高效存管。我的策略是原始数据保留 30 天用于排查问题、精确查询再通过连续查询把原始数据聚合成 1 分钟、1 小时、1 天的降采样数据长期保留用于趋势分析。-- 创建保留策略原始数据保留30天 CREATE RETENTION POLICY rp_30d ON iot_db DURATION 30d REPLICATION 1 DEFAULT -- 创建连续查询每5分钟聚合一次1分钟的均值 CREATE CONTINUOUS QUERY cq_1m ON iot_db BEGIN SELECT mean(temperature) AS mean_temperature, mean(humidity) AS mean_humidity INTO rp_30d.downsampled_1m FROM telemetry GROUP BY time(1m), deviceId, productKey END这套组合的收益很直接90 天后的数据总量里归档的降采样数据只占原始数据的几十分之一但趋势查询的性能完全不受影响。至于到底保留多久取决于业务需求——出现质量问题需要回溯的保留期长一点只是做实时监控告警的保留几天就够。6. Part 1 实测效果与已经踩过的坑到这里平台最核心的三个地基已经落地了。6 部分内容分别对应了架构决策、接入网关、设备影子、数据链路。实际跑起来后一些数值和问题也想分享出来作为 Part 1 的收官。6.1 压测数据与真实表现我用模拟器模拟了 5000 台设备同时接入 EMQX每台设备每 10 秒上报一条消息持续压测了 8 个小时。最终的数据是消息吞吐稳定在每秒 500 条左右EMQX 节点的 CPU 占用率保持在 40% 以下Redis 影子读取平均延迟 0.8ms写入平均延迟 1.2msInfluxDB 批量写入吞吐达到每秒 3000 个 point完全有余量。这个数据说明这套架构在小规模部署下完全没有压力真正的瓶颈大概率会出现在网络带宽和规则引擎的复杂处理逻辑上而不是在接入层和存储层。6.2 几个值得注意的实际问题第一个坑是 EMQX 的 CleanSession 设置。设备端如果设置了 CleanSessionfalse离线重连后会自动补发离线期间 QoS 1 的消息。这个设计本意是好但如果设备休眠时间比较长重连瞬间会一下收到大量积压消息导致设备端处理不过来崩掉。解决方法是设备端根据业务场景合理设置 CleanSession 和会话过期时间或者平台在设备上线时执行一次缓存清理。第二个坑是 InfluxDB 的字段类型冲突。同一个 measurement 里同一个字段名如果先写了整数后写了浮点数InfluxDB 会直接拒绝第二条数据。IoT 场景里尤其容易踩比如有的设备上报temperature: 25有的上报temperature: 25.6。解决方案是尽量在物模型层面把所有数值字段统一成浮点数类型或者清洗服务里做一次类型转换。第三个坑是设备上报频率不固定导致的降采样失真。有的设备只在温度变化时才上报有的设备固定 10 秒上报一次。跨设备做均值对比时如果不做加权处理统计结果会偏差很大。目前我在 Part 1 里没有做复杂处理但 Part 2 做数据分析时这会是一个必须解决的课题。6.3 留给 Part 2 的事Part 1 完成了平台的地基但距离一个真正能用的 IoT 平台还有很多路要走。我自己规划的下一个里程碑是规则引擎用户可以通过配置简单的如果温度大于 30 度就发送告警通知这种条件规则不需要写一行代码。此外告警通知的通道整合短信、邮件、App 推送、多租户权限体系、设备可视化大屏也都需要在后续版本中逐步加入。数据结构方面我已经在 Part 1 的物模型和存储设计上为这些扩展留好了接口后面填充业务逻辑时会顺利很多。最后分享一点个人体会自研 IoT 平台是一件磨刀不误砍柴工的事。前期看起来比直接用现成方案多花了很多时间但当你面对一脸这个功能平台不支持的尴尬时就会明白所有投入都是值得的。Part 1 把地基打牢了后面的功能迭代才能真正快起来。
返回列表