網(wǎng)傳感器數(shù)據(jù)全鏈路:存儲、清洗與分析實戰(zhàn)指南)
這兩年被問得最多的一個問題不是“Hadoop怎么學(xué)”而是“我們那一批傳感器每天上報幾百億條數(shù)據(jù)寫入沒問題但想查點什么一個查詢跑半小時怎么辦”。物聯(lián)網(wǎng)的傳感器數(shù)據(jù)跟傳統(tǒng)互聯(lián)網(wǎng)日志完全不是一回事——設(shè)備數(shù)量動輒百萬級每臺設(shè)備幾秒一條數(shù)據(jù)一天下來就是幾百億條記錄而且數(shù)據(jù)格式五花八門時序特征極強還有大量噪音和缺失。很多人一開始用傳統(tǒng)數(shù)據(jù)庫扛扛到幾千萬條就明顯卡頓換時序數(shù)據(jù)庫能解決一部分但涉及復(fù)雜分析、歷史歸檔、跟業(yè)務(wù)數(shù)據(jù)做關(guān)聯(lián)又力不從心。這時候Hadoop生態(tài)的價值就體現(xiàn)出來了——它不是為了存數(shù)據(jù)而存數(shù)據(jù)而是為了讓你在幾十億條傳感器記錄上還能跑出個結(jié)果來。這篇文章我從實際項目的角度把“Hadoop物聯(lián)網(wǎng)傳感器數(shù)據(jù)”這個組合拆開講數(shù)據(jù)鏈路怎么搭、清洗策略怎么做、存儲模型怎么建、跑分析時有哪些坑最后用一個典型場景復(fù)盤收尾。適合正在做物聯(lián)網(wǎng)平臺、準(zhǔn)備上Hadoop處理設(shè)備數(shù)據(jù)的團隊也適合畢業(yè)設(shè)計選了物聯(lián)網(wǎng)方向的同學(xué)們參考。1. 物聯(lián)網(wǎng)傳感器數(shù)據(jù)的四種“脾氣”為什么非Hadoop不可搞過物聯(lián)網(wǎng)的人都有體會傳感器數(shù)據(jù)跟人工錄入的“業(yè)務(wù)數(shù)據(jù)”完全兩個物種。我經(jīng)常打比方傳感器數(shù)據(jù)像是流水線上的零件——每個單獨看起來都差不多數(shù)量卻大到嚇人而業(yè)務(wù)數(shù)據(jù)像是檔案室里的文件——數(shù)量少但每份都很重要格式也要精雕細(xì)琢。拿處理文件的方式去處理零件必然出問題。1.1 高頻寫入幾秒鐘一條量級跟日志沒法比傳統(tǒng)互聯(lián)網(wǎng)日志寫入峰值一臺服務(wù)器每秒鐘幾百條就算高了但一臺工業(yè)網(wǎng)關(guān)后面掛了上百個傳感器每個傳感器3秒上報一次這個網(wǎng)關(guān)每秒就有幾十條數(shù)據(jù)。一個中型工廠幾百臺網(wǎng)關(guān)就是每秒上萬條寫入。這還不算共享單車、智慧路燈、環(huán)境監(jiān)測這類全國性場景。我見過一個項目40萬個設(shè)備每2秒一條心跳數(shù)據(jù)光是一天的數(shù)據(jù)量就是17億條2TB壓縮后。這是物聯(lián)網(wǎng)數(shù)據(jù)跟普通日志最本質(zhì)的區(qū)別——持續(xù)不斷從不睡覺全年無休。這種高頻寫入場景下傳統(tǒng)關(guān)系型數(shù)據(jù)庫的瓶頸很明顯每一行插入都要走索引、走日志寫入吞吐上不去。Hadoop生態(tài)里的HDFS解決了存儲層的大規(guī)模吞吐問題靠的是大塊順序?qū)懚鳮afka這種消息中間件則解決了“接入層削峰”的問題——后面細(xì)說。1.2 時序特征極強每條數(shù)據(jù)都帶著時間戳但時間戳最不可信傳感器數(shù)據(jù)本質(zhì)上是時序數(shù)據(jù)一個設(shè)備ID 一個時間戳 一個或多個測量值。這個特性決定了存儲模型可以高度簡化——按設(shè)備分桶按時間排序所有查詢要么是按設(shè)備查一段歷史要么是按時間范圍跨設(shè)備掃描。跟業(yè)務(wù)數(shù)據(jù)那種多表關(guān)聯(lián)的復(fù)雜結(jié)構(gòu)比起來傳感器數(shù)據(jù)的模型簡單得多但數(shù)據(jù)量卻大幾個量級。但這里有個反直覺的坑時間戳看起來是數(shù)據(jù)自帶的屬性實際上恰恰是傳感器數(shù)據(jù)里最不靠譜的字段。設(shè)備時鐘漂移、網(wǎng)關(guān)緩存重傳、網(wǎng)絡(luò)延遲都會讓數(shù)據(jù)到達(dá)時間和數(shù)據(jù)產(chǎn)生時間出現(xiàn)偏差有的甚至差幾個小時。我在后面有一節(jié)專門講這個問題這里先提個醒——如果你把入庫時間當(dāng)成了傳感器產(chǎn)生時間那后面的分析結(jié)果會很離譜。1.3 格式參差不齊幾十種設(shè)備幾十種協(xié)議聚在一起就是爛攤子一個項目里很少只有一種傳感器。溫度傳感器、濕度傳感器、振動傳感器、能耗表計、定位追蹤器……每種的報文格式都不一樣。有的上報字段叫temp有的叫temperature有的干脆是data: {v: 23.5}這種嵌在JSON里的。更過分的是不同批次的固件版本字段含義還會變。這不是代碼規(guī)范問題而是設(shè)備廠商太多、協(xié)議標(biāo)準(zhǔn)跟不上的現(xiàn)實。這部分臟活累活在Hadoop架構(gòu)里通常拆成兩層解決接入層做“格式歸一化”把亂七八糟的報文解析成統(tǒng)一的JSON或Avro格式分析層再做“字段標(biāo)準(zhǔn)化”把歷史數(shù)據(jù)統(tǒng)一到一個口徑。這也是為什么我在下一節(jié)強調(diào)Kafka的schema管理能力。1.4 數(shù)據(jù)的價值密度低但分析需求卻很高單條傳感器數(shù)據(jù)的價值密度極低——“設(shè)備A在10:03:27時刻的溫度是26.1℃”——這條信息基本沒用。但成百上千設(shè)備的時間序列拼在一起就能看出設(shè)備是否異常、產(chǎn)線是否過載、能源消耗是否異常。也就是說物聯(lián)網(wǎng)數(shù)據(jù)的特點是單體沒用聚合才有價值。這個特性決定了存儲不能丟但又不能粒度太細(xì)地長期全保留分析既要能全量掃描比如找出上個月所有設(shè)備的溫升曲線又要能快速定位比如查某個設(shè)備3小時前的瞬時值。Hadoop生態(tài)可以同時滿足這兩類需求HDFS/HBase管存儲Hive/Spark管全量分析HBase或Redis管點查加速。這也是為什么我不建議只上一個時序數(shù)據(jù)庫——時序庫確實在寫入和點查上有優(yōu)勢但跨設(shè)備復(fù)雜聚合分析、跟其他系統(tǒng)的數(shù)據(jù)做關(guān)聯(lián)還是Hadoop生態(tài)更順手。2. 整條數(shù)據(jù)鏈路怎么搭傳感器、網(wǎng)關(guān)、Kafka到HDFS先給一個我自己慣用的參考架構(gòu)再挨個拆解每一層為什么要這么選。傳感器設(shè)備 → 邊緣網(wǎng)關(guān) → Kafka數(shù)據(jù)接入層→ 流處理/清洗 → HDFS/HBase數(shù)據(jù)存儲層→ Hive/Spark分析層→ 應(yīng)用2.1 為什么中間非要加一層Kafka入湖和入庫是兩件事很多第一次做物聯(lián)網(wǎng)數(shù)據(jù)平臺的同學(xué)會問傳感器數(shù)據(jù)直接寫到HDFS不就行了干嘛要在中間加一個Kafka答案是——HDFS適合“批量落盤”不適合“每秒鐘幾萬條實時寫入”。HDFS的優(yōu)勢是大塊順序?qū)?、高吞吐批量?dǎo)入每來一條就寫入一次會產(chǎn)生大量小文件后面專門講。而Kafka的作用就是緩沖和削峰傳感器數(shù)據(jù)先沖到Kafka里下游不管是用Flume還是用Spark Streaming按自己的節(jié)奏批量寫入HDFS。這樣做還有另一層好處數(shù)據(jù)入湖和數(shù)據(jù)處理解耦了。設(shè)備不用關(guān)心下游存儲系統(tǒng)的死活Kafka里的數(shù)據(jù)可以先攢著哪怕下游HDFS集群重啟、跑批任務(wù)掛了數(shù)據(jù)一條不丟。Kafka默認(rèn)保留策略是7天這7天就是你的“后悔藥窗口”和“追數(shù)窗口”。2.2 選型對比Flume還是Kafka Connector有了Kafka之后從Kafka到HDFS這一段有三個方案經(jīng)常被拿來比Flume、Kafka Connect特別是HDFS Sink Connector、以及直接用Spark Streaming寫。我列個表用實際項目經(jīng)驗說話方案優(yōu)點缺點適合場景Flume Kafka Source穩(wěn)定上手快整套架構(gòu)都是Apache系配置繁瑣自定義Interceptor要寫Java監(jiān)控能力一般日志型數(shù)據(jù)、簡單管道Kafka Connect HDFS Sink連接器生態(tài)好支持Avro/Parquet/ORC自動分區(qū)有schema管理依賴Schema Registry版本匹配坑多兼容性問題在CDH/HDP之間尤其明顯標(biāo)準(zhǔn)化程度高、字段變更少的管道Spark Structured Streaming一步到位邊寫邊清洗結(jié)局大白于天下寫完每批數(shù)據(jù)即可做處理處理邏輯寫不好容易拖垮寫入性能資源消耗比前兩個高既要做清洗又要寫庫管道邏輯復(fù)雜我自己的偏好是管道簡單就用Flume管道邏輯復(fù)雜就直接Spark Structured StreamingKafka Connect反而用得少——因為它把很多處理邏輯限制在了配置層一旦遇到業(yè)務(wù)字段映射這類需求配置比寫代碼還痛苦。不過這是個人喜好團隊技術(shù)棧不同選擇不同有一點是公認(rèn)的從Kafka到HDFS的這層管道一定要支持按批提交、失敗重試和流量監(jiān)控缺一個后面都會很被動。2.3 存儲層選擇HDFS還是HBase兩條腿走路物聯(lián)網(wǎng)傳感器數(shù)據(jù)的存儲我建議兩條腿走路一份放HDFS一份放HBase或者用其他列式存儲。HDFS Parquet/ORC用于歷史歸檔和批量分析。按時間分區(qū)存儲保留周期可以很長一年甚至幾年。分析任務(wù)是讀這類數(shù)據(jù)。HBase用于最近N天的點查和實時查詢。比如“查某個設(shè)備當(dāng)前狀態(tài)”、“查某個設(shè)備最近一小時曲線”走HBase的RowKey索引非??臁owKey設(shè)計一般是設(shè)備ID逆序 時間戳避免熱點后面細(xì)說。如果你不想引入HBase也可以用HDFS Hive直接扛點查——建好分區(qū)表用Spark SQL按分區(qū)過濾查千萬級設(shè)備量下響應(yīng)一般在秒級到十秒級很多場景夠用了。非要說什么時候必須上HBase那只有“查詢要求毫秒到百毫秒級”的場景比如實時告警聯(lián)動、大屏點查。3. 讓數(shù)據(jù)“干凈”地入庫傳感器數(shù)據(jù)的清洗策略與實操“臟數(shù)據(jù)進(jìn)庫分析結(jié)果就是垃圾。”這句廢話在物聯(lián)網(wǎng)領(lǐng)域尤其重要——因為這行業(yè)的數(shù)據(jù)臟法跟別處不一樣不是人為錄入錯誤而是設(shè)備層面的物理性失真。采集環(huán)節(jié)不可控所以清洗策略必須在入庫前做扎實。3.1 三類最常見的傳感器“臟數(shù)據(jù)”及判據(jù)我歸納下來傳感器數(shù)據(jù)清洗主要解決三類問題1重復(fù)數(shù)據(jù)同一時刻同一設(shè)備上報了兩條一模一樣的記錄。原因一般是設(shè)備的重傳機制網(wǎng)絡(luò)抖動導(dǎo)致ACK沒到設(shè)備重發(fā)或者網(wǎng)關(guān)轉(zhuǎn)發(fā)了兩次。判據(jù)就是設(shè)備ID 來源時間戳這兩列組合去重。特別注意去重不能只看JSON是否完全一樣因為兩次重傳的數(shù)據(jù)可能部分字段不同比如接收時間不同要用業(yè)務(wù)主鍵去重。2離群值/超范圍值溫度傳感器報了500℃濕度報了-20%這種數(shù)據(jù)明顯不物理。判定方法是給每個測點配置合理的上下限。但這里有個容易踩的坑——不同應(yīng)用場景的閾值不一樣。同一塊溫度傳感器用在常溫廠房0~40℃和用在冷鏈-30~10℃上正常范圍完全不一樣。所以閾值配置必須跟著設(shè)備類型走不能寫死在處理代碼里。3缺失數(shù)據(jù)與虛假數(shù)據(jù)設(shè)備掉線會導(dǎo)致一段時間完全沒有數(shù)據(jù)設(shè)備故障則可能反復(fù)上報同一個值比如一直報24.0。缺失數(shù)據(jù)還好辦補一個空或標(biāo)記缺失即可虛假數(shù)據(jù)最難搞要結(jié)合“數(shù)值隨時間是否變化”來判定。我在項目里用過最簡單的判據(jù)同一測點連續(xù)10條數(shù)據(jù)值完全一樣就標(biāo)記為“疑似異常值”寫入清洗表里人工抽檢。3.2 清洗在哪個環(huán)節(jié)做端側(cè)、接入層、還是分析層三處都有活干但職責(zé)不同端側(cè)邊緣網(wǎng)關(guān)做格式解析、協(xié)議轉(zhuǎn)換、基礎(chǔ)校驗必填字段是否缺失、報文是否合法。這里不做復(fù)雜邏輯因為網(wǎng)關(guān)算力有限而且升級困難——盡量少給端側(cè)加戲。接入層清洗任務(wù)/流處理做去重、時間口徑統(tǒng)一、字段標(biāo)準(zhǔn)化、值域校驗。這一層是清洗的主戰(zhàn)場因為數(shù)據(jù)到了這里才被集中看到可以做跨設(shè)備或按設(shè)備類型的規(guī)則判斷。分析層做深度清洗和異常檢測。比如后面要訓(xùn)練模型或做設(shè)備健康度評估這時才做滑動窗口、趨勢判斷等復(fù)雜邏輯。一個具體的實操建議在接入層做“標(biāo)準(zhǔn)時間”字段。設(shè)備上報的device_time設(shè)備本地時間和receive_time網(wǎng)關(guān)/平臺接收時間都要保留但下游統(tǒng)一用event_time作為事件時間等于在清洗時就明確時間口徑。我在清洗任務(wù)里一般這樣規(guī)定優(yōu)先用設(shè)備本地時間但如果設(shè)備時鐘偏差跟接收時間比超過10分鐘則標(biāo)記為時鐘漂移數(shù)據(jù)改用接收時間并將原始時間放device_raw_time字段備查。后面講Spark Structured Streaming時再說watermark怎么配合這個口徑。3.3 清洗SQL長什么樣一段可以直接抄的示例假設(shè)清洗后統(tǒng)一輸出到Hive的ODS層表ods_sensor_data上游Kafka里的原始數(shù)據(jù)是JSON字符串我用Spark Structured Streaming做實時清洗核心邏輯如下省去環(huán)境初始化和參數(shù)配置只看清洗主體// Spark Structured Streaming 消費 Kafka清洗后寫入HDFS val raw spark .readStream .format(kafka) .option(kafka.bootstrap.servers, kfk01:9092,kfk02:9092) .option(subscribe, sensor_raw) .load() val parsed raw .selectExpr(CAST(value AS STRING) as json_str) .select(from_json($json_str, sensorSchema).as(data)) .select( $data.device_id.as(device_id), $data.temp.as(temp_raw), // 時間口徑統(tǒng)一設(shè)備時間優(yōu)先漂移則用接收時間 when( abs(unix_timestamp($data.sensor_time) - unix_timestamp($data.receive_time)) 600, $data.sensor_time ).otherwise($data.receive_time).cast(timestamp).as(event_time), // 值域校驗 when($data.temp.between(-40, 85), $data.temp).otherwise(lit(null)).as(temp), // 去重 $data.msg_id.as(dedup_key) ) // 以 msg_id 為主鍵做去重用stateful操作這一段可以直接做骨架去擴展。有幾點值得說明from_json解析時一定要定義好sensorSchema字段變更時要做好兼容否則一個新設(shè)備類型上來就全管道崩了。去重用msg_id而不是設(shè)備ID時間戳拼接是因為不同批次/不同協(xié)議下msg_id的生成規(guī)則可能不同。最穩(wěn)妥的做法是設(shè)備ID傳感器時間戳隨機數(shù)三段拼一個唯一ID。對于斷言失敗的異常值我沒直接丟棄而是置NULL——這樣后面做分析時能區(qū)分“沒數(shù)據(jù)”和“數(shù)值不合法”兩種含義不同。4. 查得快才算數(shù)Hive分區(qū)建模與Spark分析實踐數(shù)據(jù)入庫只是第一步。真正的價值在“查”——但很多人發(fā)現(xiàn)數(shù)據(jù)是存進(jìn)去了Hive表也建了跑一個統(tǒng)計查詢要半小時Spark任務(wù)動不動OOM。這大概率是數(shù)據(jù)模型設(shè)計出了問。題傳感器數(shù)據(jù)的查詢模式高度固定設(shè)計好了90%的分析查詢都能走分區(qū)裁剪速度能差幾十倍。4.1 分區(qū)策略按時間分區(qū)還是按設(shè)備分組我的建議傳感器數(shù)據(jù)查詢有兩個天然維度設(shè)備維度和時間維度。Hive表怎么做分區(qū)直接決定了查詢效率。**按時間分區(qū)如按小時/天分區(qū)**是默認(rèn)方案絕大多數(shù)場景都適用。原因很簡單傳感器數(shù)據(jù)分析里跨設(shè)備的時間范圍掃描是最常見的查詢——比如“查全天所有設(shè)備的平均溫度”、“查最近7天某型號設(shè)備的異常率”都是時間維度主導(dǎo)。按時間分區(qū)后這類查詢只需讀取對應(yīng)分區(qū)的數(shù)據(jù)掃描量從“全表”降到“一天的量”。但這里有幾個實操層面的建議分區(qū)粒度不要太小。監(jiān)控數(shù)據(jù)量不大的場景按天分區(qū)夠了量太大再考慮按小時。我見過有人按5分鐘分區(qū)結(jié)果一個查詢要合并幾千個分區(qū)文件MapReduce的啟動開銷比實際計算還大得不償失。分區(qū)列不要用dt這樣沒意義的字段干脆就叫event_date直接用清洗后的event_time來分區(qū)。這樣查詢時WHERE event_date 2024-06-01優(yōu)化器能精確裁剪。如果單體設(shè)備數(shù)據(jù)量極大比如一臺設(shè)備每天上千萬條記錄可以考慮“設(shè)備ID哈希分桶按時間分區(qū)”的雙層結(jié)構(gòu)。分桶字段是設(shè)備ID查詢某個設(shè)備的完整歷史時就能跳過大量不相關(guān)文件。但注意分桶數(shù)不能亂設(shè)要與文件大小匹配否則小文件問題會變本加厲。下面是建表模板可以直接抄CREATE TABLE dwd_sensor_data ( device_id STRING, device_type STRING, event_time TIMESTAMP, temp DOUBLE, humidity DOUBLE, vibration DOUBLE, ... ) PARTITIONED BY (event_date STRING) STORED AS PARQUET TBLPROPERTIES (parquet.compressionSNAPPY);4.2 列式存儲壓縮讓單條記錄再瘦一圈同樣的數(shù)據(jù)用TEXT存和用ParquetSnappy存查詢性能可以差5倍以上存儲空間可以壓縮60%以上。原理不復(fù)雜列式存儲只在讀取查詢涉及到的列時讀取對應(yīng)數(shù)據(jù)塊而傳感器數(shù)據(jù)一張表動輒幾十個字段多數(shù)查詢只用其中兩三個字段——行式存儲要把一整行讀完才能拿到一個列的值。這就像你從一疊定制的紙質(zhì)表格里查所有人的手機號行式存儲要求翻完每一張完整表格列式存儲直接把“手機號”那一列抽出來。注意Parquet的另一個好處是內(nèi)置schema列名、列類型用Hive/Spark讀時不用再指定分隔符和字段順序少了不少解析錯誤。我用的是Snappy壓縮——壓縮率比Gzip差一點但解壓速度快適合查詢頻繁的場景。冷數(shù)據(jù)想壓得更狠可以直接換ORCZlib但ORC在Spark里的支持沒Parquet那么順滑要看你的分析引擎主要用什么。4.3 Spark讀Hive數(shù)據(jù)兩個最常見的性能殺手講個真實數(shù)據(jù)一臺Spark任務(wù)讀1TB的Hive表做設(shè)備聚合分析第一次跑了一個多小時第二次五個小時第三次直接OOM。根因兩個殺手一讀出來的寬表。原始表有30個字段但分析只需要device_id、event_time、temp三個字段。代碼里如果有人用了SELECT *或者Spark的“謂詞下推”沒生效整表都被讀進(jìn)來了。解決辦法分析SQL里顯式寫出需要的列不要圖省事寫*檢查Spark物理計劃里PushedFilters是否生效。殺手二不合理的join策略。用Hive做設(shè)備基礎(chǔ)信息表和傳感器數(shù)據(jù)表的關(guān)聯(lián)如果基礎(chǔ)表只有幾萬行而傳感器表有幾十億行默認(rèn)的Shuffle Join會把幾十億行全部shuffle到所有節(jié)點傳輸量巨大。正確做法是廣播小表-- 使用Broadcast Join提示避免大表Shuffle SELECT /* BROADCAST(dim) */ s.device_id, d.region_name, AVG(s.temp) AS avg_temp FROM dwd_sensor_data s JOIN dim_device d ON s.device_id d.device_id WHERE s.event_date 2024-06-01 GROUP BY s.device_id, d.region_name;/* BROADCAST(dim) */這個提示能強制Spark把dim_device分發(fā)給每個Executor傳感器大表在本地完成關(guān)聯(lián)省掉一次幾億行的Shuffle。我見過很多團隊優(yōu)化半天沒效果最后就是加了這個提示瞬間提升性能。4.4 實時分析怎么做Structured Streaming與“延遲數(shù)據(jù)”處理物聯(lián)網(wǎng)場景里實時和準(zhǔn)實時是一對繞不開的需求。要么是“設(shè)備數(shù)據(jù)延遲多久能看到”要么是“告警規(guī)則能不能在秒級觸發(fā)”。我的經(jīng)驗是絕大部分物聯(lián)網(wǎng)場景不需要真正的毫秒級實時流處理秒級到分鐘級的準(zhǔn)實時就夠了。我也是這么落地的Kafka里取數(shù)據(jù)Spark Structured Streaming每30秒觸發(fā)一次micro batch做清洗和簡單聚合后寫入結(jié)果表。關(guān)鍵在哪延遲數(shù)據(jù)。設(shè)備掉線一段時間后重新上線會把歷史緩存數(shù)據(jù)一股腦傳上來導(dǎo)致流任務(wù)里出現(xiàn)“昨天的事件今天才到達(dá)”。處理不好聚合結(jié)果會來回跳看板上的數(shù)字忽高忽低。解決方法是watermark水印機制——告訴流引擎“允許遲到多久”超時的一律丟棄或單獨走補償流程。我一般設(shè)10分鐘的watermark跟前面清洗時“設(shè)備時鐘漂移超過10分鐘改用接收時間”的口徑保持一致// 水印機制處理延遲數(shù)據(jù) events .withWatermark(event_time, 10 minutes) .groupBy(window($event_time, 1 minute), $device_id) .agg(avg($temp).as(avg_temp))5. 上線后才會遇到的三個經(jīng)典坑時鐘漂移、小文件與寫入熱點這一節(jié)寫的都是我在生產(chǎn)環(huán)境里真實踩過的坑踩一次抖三抖的那種。前兩個講了理論基礎(chǔ)這里專門講故障現(xiàn)場和修復(fù)過程。5.1 坑一設(shè)備時鐘漂移把“峰值分析”做成了“災(zāi)難現(xiàn)場”有個項目做工廠電力負(fù)荷分析目標(biāo)是看設(shè)備集群在“哪個時間段”用電最猛。上線兩周后BI團隊反饋說數(shù)據(jù)完全沒法看——凌晨3點出現(xiàn)用電高峰白天反而波谷這跟工廠作息完全不符。排查過程是這樣的先看Kafka里的原始數(shù)據(jù)設(shè)備時間戳是正常的白天8點再查Hive表發(fā)現(xiàn)event_time字段竟然變成了凌晨3點。問題出在清洗任務(wù)——我用的是unix_timestamp($data.sensor_time) - unix_timestamp($data.receive_time)來判斷時鐘偏差大于10分鐘就改用接收時間。但有個批次的網(wǎng)關(guān)固件有Bug每次重啟后本地時鐘會回退8小時。這些設(shè)備上報的sensor_time比服務(wù)器的receive_time晚8小時——注意是晚數(shù)值小絕對值差剛好480分鐘。我的代碼里判斷條件是絕對值大于600秒只能識別“設(shè)備時間超前”識別不了“設(shè)備時間落后8小時”這種情況。結(jié)果這批設(shè)備的所有數(shù)據(jù)都被當(dāng)成漂移數(shù)據(jù)處理用接收時間替換了傳感器時間可接收時間卻是服務(wù)器收到的時刻——凌晨3點。修復(fù)思路把“基于單條記錄的時鐘漂移判斷”改成“基于設(shè)備維度的連續(xù)漂移監(jiān)控”。對每臺設(shè)備持續(xù)統(tǒng)計receive_time - sensor_time的差值分布如果這個差值在一段時間內(nèi)穩(wěn)定在一個非0值附近就說明設(shè)備時鐘存在固定偏移應(yīng)該按補償量修正而不是直接丟棄。另外最終判斷永遠(yuǎn)以業(yè)務(wù)上“用電高峰在白天”這條規(guī)則校驗數(shù)據(jù)是否符合正常模式——這類簡單的業(yè)務(wù)合理性檢查往往能最快發(fā)現(xiàn)問題。5.2 坑二KafkaFlume寫入HDFS導(dǎo)致的小文件堆成山小文件問題是Hadoop環(huán)境里的經(jīng)典殺手。當(dāng)時用的Flume從Kafka拉數(shù)據(jù)默認(rèn)每500條提交一次每提交一次就到HDFS寫一個文件——數(shù)據(jù)量一大一天就產(chǎn)生幾萬個“小文件”。HDFS的NameNode每個文件大約占150字節(jié)元數(shù)據(jù)內(nèi)存幾千萬個文件就能吃掉幾個GB內(nèi)存更重要的是Spark/Hive跑分析時要列出和處理幾十萬個文件光打開文件的時間就比計算時間長。我定位到的根因是兩層Flume的batchSize設(shè)得太小500條hdfsSink的分區(qū)策略又按時間分了太細(xì)。解決過程用了三板斧加大Flume的batchSize到1000~5000讓每個批次攢更多數(shù)據(jù)再提交。調(diào)整hdfsSink的rollInterval和rollSize——不要讓文件每幾分鐘就滾動一次設(shè)置成“文件超過128MB或30分鐘才滾一次”充分利用大塊寫入保證吞吐。上完這兩步還沒根治后來加了一層HBase做緩沖層數(shù)據(jù)先寫入HBaseHBase再定期合并Compaction后輸出HFile到HDFS完全繞開“Kafka到HDFS直寫”的小文件問題。現(xiàn)在這個項目中我們最推薦的做法是Kafka → Spark Streaming → HDFS的方式寫入時直接用coalesce控制輸出分區(qū)數(shù)強制生成足夠大的文件。同一批數(shù)據(jù)不控制分區(qū)數(shù)可能生成幾百個小文件coalesce(4)直接把輸出收斂為4個大文件。問題看上去是“文件多”本質(zhì)是“提交粒度太碎”理解了這一點就明白各種方案的本質(zhì)都指向同一個方向——提高單文件體積。5.3 坑三寫入熱點——RowKey設(shè)計失敗導(dǎo)致HBase單節(jié)點被打爆HBase處理實時數(shù)據(jù)時RowKey設(shè)計是性命攸關(guān)的事情。我們曾直接用了設(shè)備ID 時間戳作RowKey??雌饋頉]毛病——但傳感器數(shù)據(jù)的時間是持續(xù)遞增的于是所有寫入都集中在同一個Region Server上。那是單Region熱點直接導(dǎo)致某個節(jié)點負(fù)載極高其他節(jié)點閑著。修復(fù)方案是把RowKey改成設(shè)備ID逆序 時間戳。比如設(shè)備ID從device_000000000001變成100000000000_device再拼上時間戳。這樣相同設(shè)備的不同時間點在主鍵上分布到了不同的Region寫壓力同時分?jǐn)偟郊豪锏亩鄠€節(jié)點上。另一個常見備選方案是在RowKey前面加一個隨機前綴比如把device_id哈希后取前兩位但這么做的代價是查詢時不知道前綴Scan就不連貫。設(shè)備ID逆序是“分散寫”和“方便查”之間比較平衡的方案——反查時我們知道完整的設(shè)備ID照樣能直接命中RowKey前綴不會犧牲查詢性能。修復(fù)之后監(jiān)控圖上寫吞吐從“一個Region扛”變成“幾個Region平攤”P99延遲從200ms降到40ms。這個經(jīng)驗后來也驗證了設(shè)計中另一個判斷物聯(lián)網(wǎng)高吞吐場景下熱點問題本來就是第一殺手RowKey設(shè)計永遠(yuǎn)要先想寫熱點再想查便利性。6. 完整案例復(fù)盤一個百萬級環(huán)境監(jiān)測平臺從數(shù)據(jù)接入到分析落地的全過程前面講了這么多方法論和坑最后用一個我參與過的真實項目把這些串起來。這個項目不需要透露具體甲方名字就講場景和數(shù)據(jù)。項目背景幾十個城市、上萬個監(jiān)測點每個監(jiān)測點部署PM2.5、溫濕度、風(fēng)速風(fēng)向、噪聲等7種傳感器每10秒上報一條數(shù)據(jù)。全部設(shè)備累計每天產(chǎn)生約8億條數(shù)據(jù)單條原始報文是一條JSON字符串大小約300字節(jié)。這么算下來一天原始數(shù)據(jù)約24GB一年接近9TB。目標(biāo)有三實時大屏分鐘級展示各城市平均PM2.5實時告警濃度超閾值觸發(fā)預(yù)警離線分析按周/月輸出空氣質(zhì)量趨勢報告評估不同區(qū)域污染源影響歷史追溯對任意監(jiān)測點查詢?nèi)我膺^去一天的分鐘級曲線響應(yīng)5秒內(nèi)。6.1 數(shù)據(jù)流轉(zhuǎn)鏈路和核心參數(shù)采集端網(wǎng)關(guān)統(tǒng)一上報到MQTT BrokerBroker直接轉(zhuǎn)發(fā)到Kafka的sensor_raw主題。不直接讓設(shè)備連Kafka因為Kafka是TCP協(xié)議設(shè)備用MQTT連接生態(tài)更成熟。這就是協(xié)議轉(zhuǎn)換標(biāo)準(zhǔn)做法設(shè)備→MQTT平臺內(nèi)部→Kafka兩者在Broker層對接。接入清洗Spark Structured Streaming消費sensor_raw做去重、值域校驗、時鐘漂移修正并補上event_date分區(qū)字段用event_time轉(zhuǎn)換寫入HDFS的ODS層按天分區(qū)。實時部分同一個流任務(wù)同時做分鐘級聚合寫入HBase Redis緩存最近5分鐘數(shù)據(jù)供大屏點查。離線分析夜里用Spark SQL從ODS層讀全量原始數(shù)據(jù)做多維度聚合到DWD層按城市、小時、測點類型分層報表和趨勢分析直接查DWD層秒級響應(yīng)。Kafka的配置也有講究——關(guān)鍵topic分區(qū)數(shù)設(shè)為24個等于Broker數(shù)保證寫入不傾斜且消費并發(fā)可以到24。HDFS塊大小用默認(rèn)的128MB清洗任務(wù)輸出文件控制在128MB以上這個目標(biāo)所以Spark寫入時用repartition(24)。6.2 上線的時踩過的三個小問題第一個問題是MQTT到Kafka的鏈路丟數(shù)據(jù)。設(shè)備重連時網(wǎng)關(guān)會補傳斷點數(shù)據(jù)但Broker轉(zhuǎn)發(fā)到Kafka時又用了異步send缺少ACK確認(rèn)壓力大時有幾條消息被丟掉。排查后發(fā)現(xiàn)是Broker的QoS設(shè)置是“發(fā)完不管”改成QoS1至少一次和Kafka的ackall補上了丟棄缺口。第二個是Spark處理數(shù)據(jù)傾斜。某幾個工業(yè)區(qū)的監(jiān)測點數(shù)量明顯多于普通區(qū)域按區(qū)域聚合時出現(xiàn)了少數(shù)組件處理多倍數(shù)據(jù)的現(xiàn)象。解決辦法是按“城市小時”做二次worker分配——代價是多一次shuffle但從結(jié)果看這點額外消耗完全值得。第三個是HBase的Compaction風(fēng)暴。高峰期大量Region同時做Compaction導(dǎo)致某個RegionServer的IO被打滿反而降低查詢TPS。后來錯峰合并把HBase的自動Compaction關(guān)掉每天凌晨用運維腳本按RegionServer逐個手動執(zhí)行Major Compaction問題就平了。這也是運維層面一個合理的取舍——犧牲一點自動性換取全天服務(wù)穩(wěn)定。6.3 大盤數(shù)據(jù)長什么樣上線穩(wěn)定幾周后我摘過幾個關(guān)鍵數(shù)據(jù)全鏈路端到端延遲平均30秒左右傳感器到平臺可查大屏數(shù)據(jù)秒級更新離線每天刷數(shù)任務(wù)2小時完成8億條原始數(shù)據(jù)Hive查詢分鐘級的趨勢分析0.8秒到3秒存儲空間方面ODS層原始數(shù)據(jù)壓縮后4.3GB/天DWD層聚合數(shù)據(jù)僅500MB/天但已經(jīng)能滿足95%以上的查詢需求。這個壓比其實很典型——原始數(shù)據(jù)全量保存一份便宜精細(xì)分析集做一份快兩層都保住。7. 一些經(jīng)驗沉淀給后續(xù)做同類項目的人幾點忠告全項目走完一遍我個人最大的體會有三條第一“接入容易治理難”。傳感器數(shù)據(jù)的接入環(huán)節(jié)看上去每個設(shè)備都連好就行實際上治理工作要占整個項目的70%以上——而且這些工作如果沒有事先規(guī)劃等到數(shù)據(jù)上線之后再做成本和被動程度都遠(yuǎn)超想象。清洗規(guī)則、主鍵設(shè)計、時間口徑應(yīng)該在數(shù)據(jù)流向設(shè)計階段就定下來而不是上線了再返工。第二“能落到HDFS的絕不浪費在內(nèi)存”。物聯(lián)網(wǎng)數(shù)據(jù)的體量決定了內(nèi)存很貴把幾百億行數(shù)據(jù)長期放Redis或內(nèi)存數(shù)據(jù)庫財力上往往撐不住。HDFS和對象存儲都用壓縮配合是最劃算的歷史存儲方式。而內(nèi)存、SSD這些高速資源只留給最需要點查和分析的層。第三“所有問題都能在監(jiān)控里現(xiàn)原形”。我們吃了不少“數(shù)據(jù)錯了但沒人發(fā)現(xiàn)”的虧——不是任務(wù)報錯而是結(jié)果不合常識。后來給Kafka lag、HDFS文件數(shù)、Region熱點、Spark任務(wù)失敗率都做了實時監(jiān)控和告警并對“峰值溫度是否異常偏移”這類業(yè)務(wù)結(jié)果設(shè)置了規(guī)則校驗。數(shù)據(jù)平臺最怕的不是壞是壞了沒人知道。這個架構(gòu)可能不是最優(yōu)解但它是經(jīng)過了生產(chǎn)環(huán)境驗證和故障打磨的。物聯(lián)網(wǎng)數(shù)據(jù)處理沒有銀彈核心是抓住自己的場景特性高吞吐、時間序列、多源異構(gòu)、低價值密度——把存儲、清洗、分析都圍繞這四個特征去設(shè)計整體不會跑偏。