IoT数据采集全链路

发布时间:2026/8/4 20:38:26
IoT数据采集全链路 工地上有一台塔机起重机上面装了一个传感器盒子每隔几秒钟就把当前的高度、重量、风速、倾斜角度等数据发出来。问题是——这些数据怎么从工地上的铁疙瘩最终变成网页上能看到的图表这就是 IoT 数据采集链路要干的事。链路图第一步设备发数据工地上的塔机和升降机上面装了传感器和通信模块。设备每隔几秒钟就把当前的状态数据高度、重量、风速、倾斜角度等通过 4G 网络发出去。这些数据发到哪呢发到一个叫EMQX的消息中转站。第二步EMQX 消息中转站EMQX 是一个开源的 MQTT 服务器可以理解成一个快递分拣中心。所有设备发来的数据都先到这里然后按照主题分类。主题长这样crane/rt/equip/厂商/项目/用户/设备编号crane/开头 塔机数据lift/开头 升降机数据rt/equip/ 实时运行数据st/equip/ 静态参数设备型号之类的rt/operator/ 操作人员信息代码里用通配符crane/////订阅所有塔机主题意思是不管哪台塔机、哪个工地的只要是 crane 开头的我都收。第三步消息回调——按主题分发项目启动时会连接 EMQX并注册一个回调函数MqttPushCallback。每当有新消息到达EMQX 就会自动调用这个回调函数的messageArrived方法。回调函数做的事情很简单看主题里含什么关键词就调哪个处理方法。主题含 crane/rt/equip/ → 处理塔机实时数据 主题含 lift/rt/equip/ → 处理升降机实时数据 主题含 st/equip/ → 更新设备静态参数 主题含 rt/operator/ → 更新操作人员信息如果处理出错了还有个错误重推机制把失败的消息重新塞回去再试一次。第四步协议解析——把乱码变成结构化数据这里有个现实问题不同厂家的设备发出来的数据格式完全不一样。比如嘉联厂家的塔机字段叫serialNumber、height、loadWeight建安厂家的塔机同样的信息字段叫devsn、height还要乘以 10、weight还要乘以 1000因为单位不同。所以写了AnalysisTerminalDataUtil针对嘉联和建安两个厂商各写了一套解析方法。不管原始格式怎么乱解析完都统一成同一个TerminalTowerData对象。最巧妙的是报警状态的解析。嘉联升降机用一个systemStatus字段存了 16 种报警状态每一位代表一种第1位 1 → 重量预警 第2位 1 → 重量报警 第3位 1 → 高度预警 第4位 1 → 高度报警 ...代码用state 1、state 4这种位运算从一个数字里拆出 16 种报警。一个字段当好几个字段用非常省空间。第五步缓冲队列——攒够一批再入库最亮的设计这是整个链路最聪明的地方。塔机每几秒发一条数据如果每来一条就往数据库写一次假设 1000 台设备同时在线每秒就是几百条 INSERT数据库扛不住。正确的做法是数据先扔进一个队列ArrayBlockingQueue攒着攒够 limitNumber 条配置的批量大小一次性批量写入数据库。数据来了 → 放进队列 → 队列满了吗 ├ 没满 → 继续等 └ 满了 → 一次性批量写入但还有个问题如果数据量小队列一直不满数据岂不是一直卡在队列里所以加了个定时兜底每 10 分钟强制把队列里剩余的数据刷进数据库不管满没满。这样保证了数据最多延迟 10 分钟入库不会丢。第六步Redis 在线检测——不用查数据库就知道设备在不在线传统做法是定时去数据库查每个设备最后上报数据的时间超过一定时间没上报就标记为离线。设备多了这个查询很慢。更聪明的做法是每收到一条数据就在 Redis 里刷新一个 key设 600 秒过期。设备一直在发数据 → key 一直被刷新 → 一直不过期 → 在线设备坏了不发数据了 → key 600 秒后自动过期 → 离线不用写任何定时任务去轮询Redis 自己的过期机制就帮你搞定了。而且 key 过期前如果设备恢复发数据还能自动改回在线状态。第七步双数据库存储数据最终存到两个地方MySQL存实时数据 业务数据设备信息、项目信息等TDengine专门存历史运行数据为什么要用两个数据库MySQL 查最近一条数据没问题但如果你要查过去一个月某台塔机每分钟的高度变化那是几十万条记录MySQL 查起来很慢。TDengine 是专门为时序数据设计的查这种历史趋势又快又省空间。代码里用DS(tdengine)注解一行代码就能切换到 TDengine 数据源查询。