多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

智能家居数据存储方案:从分布式架构到冷热分层实践

智能家居数据存储方案:从分布式架构到冷热分层实践 先说一个挺真实的场景我接过一个智能家居平台的存储项目前期没想的那么复杂结果设备一接入数据就像水龙头没关一样涌进来。温湿度、门磁、人体感应、设备心跳、告警事件、摄像头片段……每样都是数据每样都占地方。当时用户要求很简单别丢数据别让查询卡成幻灯片成本别失控。这三个要求听起来朴素实际上把存储方案的所有关键点全逼出来了。这篇文章就把我这套“大数据领域分布式存储的智能家居数据存储”的完整思路写出来包括数据怎么分类、容量怎么估算、组件怎么选、坑怎么踩给正在做智能家居或物联网平台的朋友一个能直接上手的参考。1. 先搞清楚智能家居到底在产生什么数据很多人一上来就搜“分布式存储”“Hadoop”但连自己要存的数据长什么样都没捋明白。智能家居的数据不是一张表能装下的它至少分成四类每一类的存储要求完全不一样。1.1 四类核心数据别把鸡蛋放一个篮子第一类是设备元数据包括设备型号、固件版本、所属房间、配对时间、离线状态这类信息。它的特点是“小、少、要常读”一个家庭也就几十个设备但每次控制面板打开都要拉一遍设备列表所以查询延迟要低适合放关系型数据库或Redis缓存。第二类是时序状态数据就是各种传感器的读数——温度、湿度、PM2.5、光照度、门窗开关状态、电表功率。这类数据是最“烦人”的体量主力每个设备每隔10秒到60秒上报一次写多读少而且读的时候往往是“最近1小时的变化趋势”“昨晚3点的数值是什么”天然带时间维度。第三类是事件日志比如摄像头检测到移动、门锁被指纹打开、烟雾报警器触发这类数据偶发但极其重要安全审计就靠它不能丢也不能乱。第四类是媒体数据主要是摄像头抓拍的图片和视频片段体积大、访问频次低属于典型的冷数据。把这四类混在一个存储里是最常见的错误。我见过有人把所有设备消息全放KafkaKafka一挂全完蛋也有人把视频直接塞MySQL用不了几个月数据库就膨胀到几个T查询直接超时。正确的思路是分开对待元数据走MySQL/Redis时序数据走专门的时序数据库日志走Kafka分布式存储离线归档媒体数据直接扔对象存储或者HDFS冷目录。1.2 数据量级估算从一户到一栋楼做存储方案之前必须先算容量否则后面全在空谈。我给你一个可以直接套用的模型。假设一个普通家庭有50个传感器设备每个设备平均每30秒上报一条状态单条约200字节。一天的条数就是50 × 2880 × 1.4 20万条多算了一点告警和波动一条200字节一天是40MB一年约15GB。时序数据在分布式存储里即使压缩得好也要预留2到3倍的中间空间所以一个家庭一年约50GB的原始副本占用。但智能家居从来不是按“一户”计的。一套系统跑的是整栋楼、整个小区一个社区200户那一年就是50GB × 200 10TB加上摄像头视频这个数字立刻翻到50TB到100TB。单机MySQL在5TB左右就开始吃力10TB以上的结构化时序数据没有分布式方案根本稳不住。还有一个隐藏点写入峰值。晚上6点到10点家人回家设备联网、传感器告警、摄像头被触发写入量可能是白天的5到10倍。做容量规划时峰值HOT区要比均值大3倍以上我一般直接按“平均写入速率×8”来设计Kafka的分区数和磁盘带宽。1.3 存储需求三句话不丢、能查、花得值把需求翻译成工程师的语言就是三句第一端到端写入不丢数据这要求从设备上报到最终落盘的每个环节都有ack确认和重试机制尤其Kafka的acksall参数不能省。第二查询能在秒级返回“最近一小时”“最近一周”的聚合结果这要求热数据必须走索引友好的存储比如时序数据库按时间设备ID组织索引而不是傻乎乎地扫全表。第三总体成本可控。智能家居设备本身利润薄存储不能比设备还贵所以冷热分层是必选项热数据放SSD冷数据放普通HDD甚至归档存储单GB成本差异大概在3到5倍。2. 分布式存储方案怎么选别被Hadoop忽悠很多教程一说到“大数据”就默认上Hadoop。但实际上HDFS的随机读写性能并不适合在线查询它对海量小文件的处理更是糟糕。智能家居数据存储中最典型的方案路径是什么我来拆一下。2.1 为什么单机数据库最先被淘汰用MySQL做过物联网项目的人都懂表一千万行后插入还能忍但带时间范围加设备ID的查询一旦没命中索引慢查询日志能刷几屏。就算做了分表按日分表一年365张表跨表聚合又麻烦。更致命的是单机磁盘容量有限我见过一个朋友用4T的机器跑家用的数据平台半年后磁盘告警只能手动清理历史数据然后用户找来要三个月前的记录拿不出来尴尬。单机不是不能用是它的扩展能力到头了而且缺乏数据副本硬盘一坏数据全没——智能家居安全日志丢了这是事故级问题。2.2 主流方案横向对比各方法各的活我把实际调研过的方案列个表别迷信某一家合适才是王道方案适合场景优势劣势成本量级MySQL/PostgreSQL 分库分表设备元数据、业务配置事务可靠、查询灵活扩容要改代码时序量大后扛不住低HDFS (Hadoop)离线批量分析、海量文件归档吞吐极大、扩展性强、生态成熟延迟高、小文件不友好、运维重中HBase海量明细的随机读写分布式扩展、列式存储运维门槛高聚合能力弱中高InfluxDB/TimescaleDB/IoTDB时序数据在线查询写入快、压缩好、聚合函数强单机版有容量墙、集群版偏贵中Kafka数据管道缓冲、消息削峰高吞吐、持久化、可回溯不适合做最终存储库查询弱中在智能家居场景我的结论是“橄榄形”在线查询用时序数据库离线归档用HDFS两者之间用Kafka和流处理做桥梁。纯Hadoop的方案只适合跑T1的报表分析比如“这个月每户平均用电多少”不适合实时的“现在屋里温度多少”。2.3 推荐架构一条能落地的数据流水线直接上我用的方案结构相当清晰。端点是设备网关它负责把各种协议MQTT、CoAP、私有TCP统一成JSON消息推到Kafka。Kafka在这里起两个作用一是削峰设备端的毛刺流量不会直接打崩后端的存储二是做数据缓冲下游无论是时序库还是HDFS都按自己的节奏消费互不拖累。然后是流处理层我用Flink做实时清洗和格式转换把JSON转成列式格式顺手把脏数据过滤掉比如温度传感器报出-300度这种明显异常再分别写入热存储和冷存储。热存储就是时序数据库保留最近7天到30天的明细数据提供API给前端大屏和App查询。冷存储是HDFS每天定时把超过保留期的数据从时序库导出成Parquet文件按日期和设备ID分区后续要用Hive或Spark查历史趋势。这套结构的好处是任何一个环节挂了数据都不会丢Kafka有副本下游恢复后会基于offset继续消费查询快热数据不走HDFS成本省冷数据压在HDFS的普通HDD上单GB成本很低。我实际跑过半年日均千万条消息存储侧运维介入的频率从“每周处理一次”降到了“每月例行检查一次”。3. 核心配置与实操细节照着抄就行架构定完后真正拉长战线的全是配置细节和小坑。这里我把集群规划、分区策略、文件格式、压缩时机、冷热迁移全讲透。3.1 集群规划多少节点、多大磁盘才算够一个小规模智能家居数据平台我建议起步是5台机器。两台做主节点跑ZooKeeper、NameNode、Kafka Controller、Flink JobManager要求16核32G以上内存加两块SSD做系统盘和元数据盘Kafka的log目录可以挂在SSD上提升写入性能。三台做数据节点跑DataNode和Kafka Broker每个节点配两块4T的HDD不做RAID靠HDFS的副本兜底。容量按前面的估算一个小区的100TB数据副本因子按2算需要200TB物理空间三台4T×2的节点是48T其实只够一个中型设备规模所以真实项目里数据节点的数量要根据设备接入量动态加这就是HDFS的好处——加一台机器容量就能线性扩展。配置上有一个常被忽略的地方NameNode堆内存。每个文件、每个目录、每个Block都会占内存默认4G堆只能支持大约200万个对象。智能家居的数据有大量小文件——如果直接落HDFS每秒产生的新文件数会迅速击穿NameNode内存。所以聚合写入必须做见3.3。3.2 数据分区与文件格式让查询和压缩都舒服HDFS上冷数据的目录组织我推荐这样设计/data/smarthome/{device_type}/{yyyyMMdd}/{device_id}/日期在最前查询单日全量就只扫描对应日期的目录查询单设备就再定位到具体子目录。Hive建表时把日期分区做成分区字段能用谓词下推把扫描范围缩到最小。这个设计配合“设备类型再分组”比如温度传感器和摄像头数据分开因为它们的量级和查询模式完全不同摄像头文件大且少温度传感器文件小且密分开可以避免小文件挤占大数据文件的读取带宽。文件格式方面不要用文本JSON直接落HDFS文件大而且SQL查询还得套Serde。我统一采用Parquet加Snappy压缩Parquet是列式存储聚合查询只读需要的列I/O减少一半以上Snappy压缩率大概在2到3倍解压速度也快。实测同样380MB的JSON数据转成ParquetSnappy后是142MB查询“一天内每户平均温度”的Spark SQL耗时从45秒降到9秒。如果你习惯Hive生态ORC也可以两者压缩率接近但Parquet的跨框架支持更好。时序数据的保留期也在这里明确热库保留30天超过30天的按天导出到HDFS导出任务放在凌晨2点避免高峰占用带宽。导出的粒度是每小时一个Parquet文件块避免生成过多小文件。3.3 写入链路实操Kafka参数、批次聚合和压缩写入链路是从Kafka开始的。几个关键参数调不好后面全堆炸。acksall——所有的ISR副本都确认才算成功丢数据容忍度为零时必须设它retries建议设5到10min.insync.replicas设为2这样即使一台Broker挂了写请求也不会悄悄失败。Kafka的topic按数据类型拆分我最少建三个telemetry存时序读数、event存告警日志、media存视频元数据。每个topic的分区数先按“目标写入速率/单分区可支撑的5MB每秒”估算100MB每秒的峰值就跑20个分区。分区多了不好少了更不好扩容分区在Kafka里是麻烦事。从Kafka到HDFS这一段如果用Flume或Kafka Connect一定要开启“批次合并”我按两个条件触发落盘文件累加到512MB或者事件时间超过10分钟哪个先到就flush一次。这里有个容易踩的坑——如果只按时间合并低峰期也会每个小时生成一堆几十KB的小文件如果只按大小合并高峰期文件倒是大了但低峰期文件要很久才能刷出来冷数据查询就有滞后。两个条件配合小文件数量可以控制在每百万条消息只产生约100个文件。3.4 冷热分层迁移别让任务自己跑晕了冷热迁移听起来简单做起来容易出乱子。我的做法是写一个定时调度的Spark批任务每天凌晨3点跑。第一步扫描时序库中30天前的数据按设备ID和时间段取出第二步调用文件写入工具生成Parquet文件上传到HDFS对应分区目录同时更新一张“归档清单”表记录哪些批次已经归档、数据的时间范围和文件路径第三步确认HDFS上的副本数量和文件大小都正常后才去删除时序库中的旧数据。整个流程中“先写后删”是底线很多丢数据事故都是因为先删了源结果导出的文件校验不过想重导都没得导。迁移还有一个隐藏工作量视频数据不走时序库我单独开了一条“媒体转存”通道设备上报的摄像头片段先进对象存储或HDFS的临时目录然后在第二天的批次里按日期汇总压缩成更大的归档文件最后统一挪到冷区。这样做避免了视频文件长期占用在线存储的宝贵空间——一段1080P的视频动辄几十MB如果放在热库存储成本立刻翻倍。4. 常见问题与排查实录系统跑起来后问题从来不会缺席。我把几个我亲历过的高频问题写下来附带排查思路和解决办法这些东西常规文档里不会写。4.1 小文件泛滥NameNode内存告急我接过一个“HDFS突然频繁Full GC”的现场一查Block数量几千万个。原因是设备类型的Kafka消息没有合并写入任务每5分钟落一次盘一天下来产生288个文件一个月近9000个半年就是5万多个每个文件对应一个BlockNameNode每个Block对象要占150字节左右的内存1000万Block就是1.5G内存很快被打爆。排查方法很简单用hdfs dfs -count统计各目录的文件数超过10万文件的目录就是重灾区再用hdfs fsck看一下Block分布。解决手段分两步第一时间把写入端改成按大小时间双条件合并见3.3避免新的小文件继续产生存量小文件用Hive的INSERT OVERWRITE重新写一遍按天合并成每个500MB到1GB的大文件。这类问题不需要重启集群但影响很大早发现早处理。4.2 数据倾斜一个分区独吞80%流量智能家居里总有“话痨设备”。有一段我观察到某个Broker的磁盘IO和网络流量明显比其他节点高Kafka的消费lag还忽大忽小。查了下topic的各个分区消息量结果发现某个分区的数据量是其他分区的15倍。原因出在分区策略上——我当时按hash(device_id) % partitions做分区而那一批设备又集中接入在一个网关下hash后全挤到同一个分区。排查时用kafka-run-class kafka.tools.GetOffsetShell直接看各分区offset差距一眼就暴露了。解决办法是调整分区键把device_id改成加上gateway_id的组合让消息散布更均匀同时把“话痨类设备”比如温度传感器群里混着的高频电量计单独拆到一个topic给它独立的分区数和消费实例避免拖累正常设备的数据读取。改完之后整个集群的吞吐立刻平稳了。4.3 时间戳漂移离线分析和实时数据的“时差”设备端的时钟大多数并不可靠。有一阵我发现HDFS上的冷数据时序图里凌晨3点那一条曲线的数据量诡异高白天的反而少。查下来是部分设备在断网重启后系统时间回到了出厂设置上报数据自带了一个晚8小时的时间戳全被归档到错误的日期分区里。这个问题在设备端很难根除所以我直接在流处理层加了校准逻辑如果设备上报时间和当前服务端时间相差超过24小时记录一条异常日志同时用服务端接收时间重新盖上时间戳对于差异小于24小时的则保留原始时间戳但打上time_corrected1的标签后续分析时完全可以选择只信服务端时间还是原始时间。这个标签字段付出的成本很低但避开了一堆“数据怎么凭空消失了一小时”的荒诞剧情。4.4 副本丢失与磁盘故障分布式不是保险箱HDFS默认副本因子是3很多新手以为这就永远不会丢数据。实际不是一只老化的HDD在坏道扩大的时候会先进入“慢盘”状态写请求全部超时DataNode会被NameNode标记为不健康它上面的副本被自动剔除如果同时段另一台机器也坏盘某个Block的两个副本就可能同时消失留下一个孤儿副本。所以我建议核心数据容量规划时就把副本因子设为3非核心冷数据设为2同时每个数据节点上保留一个空闲磁盘插槽故障时可以快速热插拔替换。监控上不要只盯着“集群健康度”要盯“Under-Replicated Blocks”这个指标任何时段的积压数量不能超过集群文件总数的0.5%一旦超过就要主动找原因。5. 最后再分享两个小技巧我踩了太多坑才意识到存储方案里最省钱也最容易被忽略的两件事拿出来收尾。第一件事监控必须覆盖端到端而不只是HDFS健康。我在生产环境里设了三个“金丝雀指标”Kafka消费lag的秒级监控、时序库写入吞吐的分钟级监控、HDFS归档任务成功率的日级监控。任何一项异常都会直接触发告警邮件和群消息。没有这些你会发现数据出了问题往往是在用户投诉一周以后那时候再倒查元数据就非常被动了。第二件事数据压缩的成本远低于扩容的成本。ParquetSnappy的组合让我的HDFS冷数据体积缩减到原始JSON的40%左右。配合冷热分层热库只占全部存储的15%其余的85%都躺在便宜的冷存储里。这一点点“存前按列式压缩”的习惯让集群的支撑设备数量直接翻倍。智能家居的数据存储没有一个放之四海而皆准的方案但“分开存、分层存、先合并再落盘”这三条经验是通用底层的逻辑。业务规模小时也不要贪多先跑起来再逐步加HDFS和Spark。就算以后数据量再涨只要写入端聚合和目录分区设计没跑偏换到更大集群也只是改个地址的事。
返回列表