據(jù)倉(cāng)庫(kù)中的數(shù)據(jù)清洗方法:分層架構(gòu)與工具實(shí)戰(zhàn))
我一直覺(jué)得很多團(tuán)隊(duì)把數(shù)據(jù)清洗這件事做小了。一說(shuō)數(shù)據(jù)清洗第一反應(yīng)就是寫(xiě)幾個(gè)SQL把空值填上、把重復(fù)行去掉。但在數(shù)據(jù)倉(cāng)庫(kù)這個(gè)場(chǎng)景里數(shù)據(jù)清洗根本沒(méi)有這么簡(jiǎn)單——它是建模的一部分是數(shù)據(jù)質(zhì)量的防線是后面所有報(bào)表、算法、決策能不能站住腳的地基。數(shù)據(jù)倉(cāng)庫(kù)里的數(shù)據(jù)清洗方法本質(zhì)上解決的是一堆雜亂無(wú)章的原始數(shù)據(jù)如何變成可信、可用、可追溯的分析底座這件事。這篇文章我從數(shù)倉(cāng)的分層架構(gòu)出發(fā)聊清楚清洗規(guī)則應(yīng)該釘在哪一層、用什么工具做、真實(shí)臟數(shù)據(jù)場(chǎng)景下怎么排查末尾附一段網(wǎng)約車訂單清洗的完整實(shí)戰(zhàn)鏈路。適合正在搭數(shù)倉(cāng)、做離線數(shù)倉(cāng)開(kāi)發(fā)、或者被數(shù)據(jù)怎么這么臟困擾的朋友收藏起來(lái)下次直接照著做。1. 數(shù)倉(cāng)里的數(shù)據(jù)清洗不是修數(shù)據(jù)是建模的一部分先做一個(gè)認(rèn)知上的正本清源。很多從業(yè)務(wù)系統(tǒng)轉(zhuǎn)過(guò)來(lái)的人會(huì)把數(shù)據(jù)清洗想成把這些臟值修好。但你仔細(xì)想業(yè)務(wù)系統(tǒng)里的數(shù)據(jù)清理和數(shù)倉(cāng)里的數(shù)據(jù)清洗面對(duì)的對(duì)象、約束和目標(biāo)都完全不同。1.1 業(yè)務(wù)庫(kù)清洗 vs 數(shù)倉(cāng)清洗的分工差異業(yè)務(wù)庫(kù)里的數(shù)據(jù)是給系統(tǒng)用的事務(wù)型數(shù)據(jù)庫(kù)講究的是當(dāng)前狀態(tài)正確、寫(xiě)入性能高。你很少會(huì)在業(yè)務(wù)庫(kù)里做大規(guī)模的歷史數(shù)據(jù)回溯修正因?yàn)槟菚?huì)影響線上交易。而數(shù)倉(cāng)里的數(shù)據(jù)是給分析用的它的特點(diǎn)是海量、多源、歷史累計(jì)而且洗完之后要能被反復(fù)讀取、追溯和重算。這意味著數(shù)倉(cāng)里的數(shù)據(jù)清洗至少要滿足三個(gè)額外要求可重放清洗邏輯必須是確定性的同一份輸入無(wú)論跑多少次輸出必須一致。不能像修線上數(shù)據(jù)那樣手工改幾條記錄就完事因?yàn)閿?shù)倉(cāng)要應(yīng)對(duì)的是TB級(jí)乃至PB級(jí)的批量加工??勺匪菝恳粭l被清洗過(guò)的數(shù)據(jù)最好能知道它原來(lái)長(zhǎng)什么樣、被什么規(guī)則變成了什么樣。這不僅僅是審計(jì)需求也是排查下游指標(biāo)異常時(shí)的救命稻草??煞謱忧逑磩?dòng)作一旦混在業(yè)務(wù)邏輯里后面維護(hù)的人會(huì)非常痛苦。所以清洗規(guī)則必須跟著數(shù)倉(cāng)的分層架構(gòu)走哪一層做什么邊界必須清晰。1.2 清洗規(guī)則的本質(zhì)對(duì)業(yè)務(wù)語(yǔ)義的數(shù)字化約束我打個(gè)比方。業(yè)務(wù)庫(kù)里的一條記錄就像一個(gè)剛跑完現(xiàn)場(chǎng)回來(lái)的銷售人員的草稿本字跡潦草、縮寫(xiě)隨意、甚至有些數(shù)字明顯抄錯(cuò)了。數(shù)倉(cāng)的清洗環(huán)節(jié)相當(dāng)于把草稿本重新謄寫(xiě)成一份標(biāo)準(zhǔn)化的臺(tái)賬——但不是你想怎么謄就怎么謄而是按一套明確的規(guī)范來(lái)謄。這套規(guī)范就是業(yè)務(wù)語(yǔ)義的數(shù)字化表達(dá)。什么算一個(gè)有效訂單訂單狀態(tài)為已完成且支付金額大于0就是一條業(yè)務(wù)語(yǔ)義規(guī)則它同時(shí)約束了狀態(tài)字段的取值范圍和金額字段的有效性。什么算一個(gè)正常的用戶注冊(cè)時(shí)間不為空、手機(jī)號(hào)為11位數(shù)字也是一條規(guī)則。所以數(shù)據(jù)清洗方法的設(shè)計(jì)第一步根本不是寫(xiě)代碼而是把業(yè)務(wù)規(guī)則顯式化。你在清洗之前得能回答這些數(shù)據(jù)里哪些字段是主鍵哪些字段的取值范圍是什么哪些字段之間的邏輯關(guān)系必須成立。規(guī)則列不出來(lái)代碼寫(xiě)得再漂亮都是在碰運(yùn)氣。2. ODS與DWD的分層清洗策略先收斂再深加工數(shù)倉(cāng)領(lǐng)域有個(gè)成熟的分層習(xí)慣ODS操作數(shù)據(jù)存儲(chǔ)層、DWD明細(xì)數(shù)據(jù)層、DWS匯總數(shù)據(jù)層、ADS應(yīng)用數(shù)據(jù)層。清洗動(dòng)作主要發(fā)生在ODS到DWD這一段但很多人會(huì)把ODS環(huán)節(jié)的收斂和DWD環(huán)節(jié)的清洗混在一起做最后導(dǎo)致規(guī)則散落、血緣混亂。我的經(jīng)驗(yàn)是ODS只做必做之事真正的清洗要釘在DWD層并且跟著維度建模的設(shè)計(jì)走。2.1 ODS層只做最小必要處理ODS層的定位是原始數(shù)據(jù)的鏡像。它存在的意義是保留數(shù)據(jù)的原始性讓上游來(lái)的數(shù)據(jù)在數(shù)倉(cāng)里有一個(gè)忠實(shí)的落地點(diǎn)。在這個(gè)層面上我堅(jiān)持只做三種處理增量/全量落地按業(yè)務(wù)系統(tǒng)同步過(guò)來(lái)的頻率做分區(qū)落地。技術(shù)字段補(bǔ)充比如etl_time抽取時(shí)間、source_system來(lái)源系統(tǒng)、data_version數(shù)據(jù)版本這些字段服務(wù)于后續(xù)的追溯和重算不改變業(yè)務(wù)語(yǔ)義。基礎(chǔ)編碼統(tǒng)一比如把不同來(lái)源的字符集統(tǒng)一成UTF-8把BOM頭去掉把不可見(jiàn)字符做trim處理。這些是技術(shù)層面的收斂不涉及業(yè)務(wù)判斷。一旦在ODS層做了業(yè)務(wù)規(guī)則的判斷比如過(guò)濾掉狀態(tài)為取消的訂單就出問(wèn)題了。因?yàn)閷?lái)排查數(shù)據(jù)問(wèn)題時(shí)你很難說(shuō)清楚ODS里的原始數(shù)據(jù)到底是上游就缺了還是被我們洗掉了。所以O(shè)DS的守則就一句話原樣接入只動(dòng)技術(shù)不動(dòng)業(yè)務(wù)。2.2 DWD層清洗規(guī)則要跟著維度和事實(shí)的設(shè)計(jì)走DWD層是數(shù)據(jù)清洗的主戰(zhàn)場(chǎng)。這一層要做的是把ODS里多個(gè)來(lái)源的數(shù)據(jù)按統(tǒng)一的業(yè)務(wù)定義組裝成明細(xì)事實(shí)表和維度表。清洗在這里不是孤立動(dòng)作而是建模過(guò)程的一部分。舉例來(lái)說(shuō)做訂單事實(shí)表時(shí)訂單狀態(tài)枚舉值必須統(tǒng)一。上游業(yè)務(wù)系統(tǒng)可能用0/1/2表示待支付/已支付/已取消另一個(gè)系統(tǒng)可能用P/PAID/CANCEL表示同一件事。DWD層必須把這些映射成統(tǒng)一的維度外鍵。做用戶維度表時(shí)性別、年齡、城市這些屬性的取值異常和缺失處理跟著SCD策略走緩慢變化維而不是簡(jiǎn)單填個(gè)未知了事。做事實(shí)表的外鍵關(guān)聯(lián)時(shí)要處理孤兒數(shù)據(jù)——比如訂單表里關(guān)聯(lián)不到用戶表的user_id。這種問(wèn)題用SQL的inner join是直接消失的但真實(shí)情況是這些訂單不能隨意丟棄得落進(jìn)專門的待確認(rèn)表里??梢园盐页S玫腄WD層清洗動(dòng)作按類別拆一下清洗類別典型動(dòng)作數(shù)倉(cāng)側(cè)的落地方式格式標(biāo)準(zhǔn)化日期統(tǒng)一成yyyy-MM-dd、金額統(tǒng)一成decimal(18,4)用CAST或regexp處理規(guī)則固化到ETL代碼里缺失值處理可推導(dǎo)的用業(yè)務(wù)邏輯推導(dǎo)不可推導(dǎo)的給業(yè)務(wù)默認(rèn)值或用未知維度替代用COALESCE或CASE WHEN值要可解釋重復(fù)數(shù)據(jù)剔除按業(yè)務(wù)主鍵去重保留最新或質(zhì)量最高的一條用ROW_NUMBER() OVER (PARTITION BY ...)邏輯矛盾修正比如支付時(shí)間早于下單時(shí)間這種矛盾記錄按規(guī)則重算或過(guò)濾并在明細(xì)表里打標(biāo)異常值鉗制超出物理上下限的數(shù)值如負(fù)的行駛里程區(qū)間判斷必要時(shí)置NULL并記錄2.3 數(shù)據(jù)質(zhì)量的六性檢查清單在DWD層設(shè)計(jì)清洗規(guī)則時(shí)我習(xí)慣用一個(gè)六維清單去自檢這六個(gè)維度分別是完整性、唯一性、準(zhǔn)確性、一致性、有效性、時(shí)效性。每次新的清洗邏輯加進(jìn)來(lái)就對(duì)著這六個(gè)詞過(guò)一遍缺哪個(gè)補(bǔ)哪個(gè)。完整性該有的字段有沒(méi)有。比如訂單表必須有下單時(shí)間和訂單號(hào)。唯一性主鍵不能重復(fù)重復(fù)了怎么處理、保留哪條。準(zhǔn)確性字段值和真實(shí)業(yè)務(wù)是否一致比如金額不能是負(fù)數(shù)。一致性同一個(gè)維度的編碼在不同表里必須統(tǒng)一。有效性字段值是否符合定義的取值范圍比如周幾只能是1到7。時(shí)效性數(shù)據(jù)是否在預(yù)期時(shí)間內(nèi)到達(dá)遲到的數(shù)據(jù)不能污染當(dāng)天的統(tǒng)計(jì)。這套清單不是給別人看的文檔而是你寫(xiě)清洗代碼時(shí)的一個(gè)心智框架。比如你正在寫(xiě)一段用戶地址清洗邏輯發(fā)現(xiàn)地址字段有的帶省市區(qū)有的只有一個(gè)市有的完全是空——這時(shí)候你如果不把六個(gè)維度過(guò)一遍很容易只想著填個(gè)空值而忘了校驗(yàn)非空地址里的省市區(qū)在行政區(qū)劃表里到底存不存在這個(gè)有效性問(wèn)題。3. 三個(gè)實(shí)戰(zhàn)工具的正確用法Hive SQL、pandas、Spark DataFrame清洗方法有了工具層面的選型也得聊清楚。數(shù)倉(cāng)里最常碰到的三個(gè)工具Hive SQL、pandas、Spark DataFrame它們各有擅長(zhǎng)的戰(zhàn)場(chǎng)用錯(cuò)了地方就會(huì)事倍功半。3.1 Hive SQL大規(guī)模批處理的基本盤數(shù)倉(cāng)里的絕大多數(shù)清洗工作最后還是落到Hive SQL上。原因很簡(jiǎn)單數(shù)據(jù)量大而且數(shù)倉(cāng)本身就是以Hive表為核心組織的。用SQL做清洗等于直接在數(shù)據(jù)所在的位置干活不用搞什么導(dǎo)出導(dǎo)入。SQL清洗的典型打法就是嵌套子查詢先做字段解析和格式標(biāo)準(zhǔn)化再做去重和過(guò)濾最后落表。拿一個(gè)比較常見(jiàn)的場(chǎng)景舉例——清洗用戶手機(jī)號(hào)字段WITH cleaned AS ( SELECT user_id, -- 去掉手機(jī)號(hào)里的空格、橫線、括號(hào)只留數(shù)字 REGEXP_REPLACE(phone, [^0-9], ) AS phone_raw, LOWER(email) AS email_lower, COALESCE(gender, unknown) AS gender_filled FROM ods_user_info WHERE dt ${bizdate} ), valid_check AS ( SELECT user_id, phone_raw, -- 合法性校驗(yàn)只保留11位且以1開(kāi)頭的號(hào)碼其余置NULL CASE WHEN phone_raw RLIKE ^1[0-9]{10}$ THEN phone_raw END AS phone_valid, email_lower, gender_filled FROM cleaned ) INSERT OVERWRITE TABLE dwd_user_info SELECT * FROM valid_check;這個(gè)模式好在哪每一步邏輯都體現(xiàn)在子查詢里后面的人看代碼能順著結(jié)構(gòu)反推你的清洗思路。同時(shí)Hive SQL里的REGEXP_REPLACE、RLIKE、ROW_NUMBER()、LATERAL VIEW這四板斧可以說(shuō)覆蓋了80%的格式清洗和去重需求。剩下20%計(jì)算特別復(fù)雜的才需要交給下面的工具。3.2 pandas小規(guī)模探索和規(guī)則原型驗(yàn)證pandas在數(shù)倉(cāng)體系里處于一個(gè)微妙的位置。它不適合處理億級(jí)數(shù)據(jù)但在清洗規(guī)則還不明確、你需要快速看數(shù)據(jù)畫(huà)像的階段pandas是效率最高的工具。我自己做數(shù)倉(cāng)開(kāi)發(fā)時(shí)遇到新的數(shù)據(jù)源永遠(yuǎn)是先拉一份抽樣數(shù)據(jù)到本地用pandas做探索性分析把清洗規(guī)則的原型先跑出來(lái)驗(yàn)證邏輯沒(méi)問(wèn)題再翻譯成Hive SQL投到生產(chǎn)環(huán)境。這個(gè)過(guò)程能幫你省大量時(shí)間因?yàn)橹苯釉贖ive上反復(fù)調(diào)試一個(gè)查詢可能就要等幾分鐘本地用pandas幾秒鐘就出結(jié)果了。pandas最常用的五個(gè)清洗操作可以記一下import pandas as pd df pd.read_csv(sample_data.csv, encodingutf-8) # 1. 去掉重復(fù)行保留第一次出現(xiàn)的一條 df df.drop_duplicates(subset[order_id], keepfirst) # 2. 缺失值處理按業(yè)務(wù)規(guī)則填充 df[pay_time] df[pay_time].fillna(1970-01-01 00:00:00) # 3. 異常值替換把不在合法區(qū)間內(nèi)的值替換掉 df.loc[df[mileage] 0, mileage] None # 4. 數(shù)據(jù)類型收斂統(tǒng)一日期格式 df[order_date] pd.to_datetime(df[create_time]).dt.date # 5. 自定義規(guī)則函數(shù)映射 df[order_status_std] df[order_status].map({0: pending, 1: paid, 2: cancelled})這個(gè)階段的核心產(chǎn)出不是清洗后的數(shù)據(jù)而是一套已驗(yàn)證過(guò)的清洗邏輯文檔。很多人跳過(guò)這一步直接上SQL結(jié)果規(guī)則在數(shù)據(jù)量大了之后才發(fā)現(xiàn)有問(wèn)題返工成本特別高。3.3 Spark DataFrame中大規(guī)模復(fù)雜清洗的折中方案當(dāng)數(shù)據(jù)量在千萬(wàn)到億級(jí)別而且清洗邏輯涉及復(fù)雜計(jì)算、多階段狀態(tài)處理時(shí)純SQL寫(xiě)起來(lái)會(huì)很別扭本地pandas又跑不動(dòng)這時(shí)Spark DataFrame是一個(gè)很好的折中。典型場(chǎng)景比如用Geohash做軌跡數(shù)據(jù)清洗、基于用戶行為序列做狀態(tài)推演、多張表關(guān)聯(lián)后做復(fù)雜的窗口計(jì)算。這幾個(gè)場(chǎng)景在Spark里寫(xiě)起來(lái)更接近編程思維比堆一大坨SQL直觀得多。from pyspark.sql import functions as F from pyspark.sql.window import Window # 按訂單分組按打點(diǎn)時(shí)間排序用于亂序軌跡修正 w Window.partitionBy(order_id).orderBy(F.col(point_time).asc()) df_cleaned df_raw.withColumn( is_duplicate, F.row_number().over(w) ).filter( F.col(is_duplicate) 1 ).withColumn( speed_kmh, F.lit(120) ).filter( F.col(distance_km) / F.greatest(F.col(interval_hour), F.lit(0.001)) F.col(speed_kmh) )Spark的調(diào)試成本比pandas高所以我的原則是先pandas出原型再Spark上生產(chǎn)。兩個(gè)工具使用同一套規(guī)則定義能最大程度減少翻譯過(guò)程中引入的偏差。工具選型小結(jié)用一張表說(shuō)人話工具數(shù)據(jù)量級(jí)最佳使用場(chǎng)景主要劣勢(shì)Hive SQL億級(jí)以上離線批處理、標(biāo)準(zhǔn)化清洗、去重過(guò)濾復(fù)雜計(jì)算寫(xiě)起來(lái)費(fèi)勁調(diào)試慢pandas百萬(wàn)級(jí)以下探索性分析、規(guī)則原型驗(yàn)證、一次性修正內(nèi)存瓶頸不適合生產(chǎn)大規(guī)模任務(wù)Spark DataFrame千萬(wàn)到十億級(jí)復(fù)雜清洗邏輯、軌跡/行為序列處理集群資源開(kāi)銷大原型階段成本高4. 一次網(wǎng)約車訂單清洗實(shí)戰(zhàn)從臟數(shù)據(jù)發(fā)現(xiàn)到規(guī)則落地的完整鏈路理論講再多不如跑一遍真實(shí)案例。這里用我之前做過(guò)的網(wǎng)約車訂單數(shù)據(jù)清洗項(xiàng)目來(lái)走一遍完整鏈路。這個(gè)項(xiàng)目的數(shù)據(jù)源包括訂單表、司機(jī)軌跡表、計(jì)價(jià)表三張ODS表接進(jìn)來(lái)以后問(wèn)題非常多。4.1 臟數(shù)據(jù)初檢先看數(shù)據(jù)畫(huà)像拿到數(shù)據(jù)第一天我不會(huì)立刻寫(xiě)清洗規(guī)則而是先做數(shù)據(jù)畫(huà)像。所謂畫(huà)像就是看每一列的空值率、去重率、枚舉值分布、最大最小值、數(shù)值范圍。這一步用pandas跑特別快。當(dāng)時(shí)發(fā)現(xiàn)的核心問(wèn)題有這么幾類訂單表的finish_time字段缺失率高達(dá)23%。這不可能是正常的反手去查上游原來(lái)是司機(jī)端APP在部分場(chǎng)景下沒(méi)有回傳完成時(shí)間。軌跡表里有大量打點(diǎn)時(shí)間早于訂單創(chuàng)建時(shí)間的記錄。也就是說(shuō)軌跡的時(shí)間戳亂序了。存在同一訂單號(hào)出現(xiàn)兩次的情況而且兩次的金額還不一樣。有個(gè)別訂單的行駛里程是負(fù)數(shù)。這些單看一條都會(huì)覺(jué)得這數(shù)據(jù)怎么回事但放在一起說(shuō)明清洗規(guī)則不能只做一個(gè)填缺失值得設(shè)計(jì)一整套基于業(yè)務(wù)約束的校驗(yàn)邏輯。4.2 異常軌跡點(diǎn)排查用物理規(guī)則過(guò)濾軌跡數(shù)據(jù)的清洗是網(wǎng)約車場(chǎng)景里最有代表性的。它的臟數(shù)據(jù)主要來(lái)自GPS漂移——某個(gè)打點(diǎn)位置突然跳到幾十公里外的另一個(gè)城市或者速度計(jì)算出來(lái)遠(yuǎn)超物理上限。處理邏輯是這樣的軌跡點(diǎn)按時(shí)間排序后計(jì)算相鄰打點(diǎn)之間的距離和時(shí)間差然后算出平均速度。如果速度超過(guò)一個(gè)物理上限比如120km/h這個(gè)點(diǎn)就很可能是漂移點(diǎn)需要剔除。這個(gè)邏輯在Hive里可以用lag窗口函數(shù)實(shí)現(xiàn)WITH trajectory_sorted AS ( SELECT order_id, lng, lat, point_time, LAG(point_time) OVER (PARTITION BY order_id ORDER BY point_time) AS prev_time, LAG(lng) OVER (PARTITION BY order_id ORDER BY point_time) AS prev_lng, LAG(lat) OVER (PARTITION BY order_id ORDER BY point_time) AS prev_lat FROM ods_trajectory ) SELECT order_id, lng, lat, point_time FROM trajectory_sorted WHERE prev_time IS NULL OR ST_DISTANCE( ST_POINT(prev_lng, prev_lat), ST_POINT(lng, lat) ) / (UNIX_TIMESTAMP(point_time) - UNIX_TIMESTAMP(prev_time)) 120這里有個(gè)細(xì)節(jié)值得說(shuō)剔除漂移點(diǎn)的時(shí)候不能只算這個(gè)點(diǎn)本身合不合理要看它和前后點(diǎn)的關(guān)系。一個(gè)點(diǎn)本身在正常城市范圍內(nèi)但和上一個(gè)點(diǎn)之間隔了300公里這就是明顯的漂移。物理速度約束比單純的范圍約束更有效。4.3 時(shí)間亂序與重復(fù)訂單的處理思路時(shí)間戳亂序在物聯(lián)網(wǎng)和APP上報(bào)場(chǎng)景里經(jīng)常出現(xiàn)。網(wǎng)約車軌跡表的打點(diǎn)時(shí)間理論上必須大于等于訂單創(chuàng)建時(shí)間、小于等于訂單完成時(shí)間但實(shí)際數(shù)據(jù)里會(huì)出現(xiàn)亂序、重復(fù)打點(diǎn)、時(shí)間超前等情況。我的處理方式是分層解決完全亂序的記錄按order_id分組用row_number按point_time重新排序亂序但不丟失修正為正確順序。重復(fù)打點(diǎn)同一秒內(nèi)重復(fù)上報(bào)的軌跡點(diǎn)保留第一個(gè)其余剔除。時(shí)間超前point_time早于訂單創(chuàng)建時(shí)間這類記錄要么是設(shè)備時(shí)鐘問(wèn)題要么是緩存上報(bào)問(wèn)題直接剔除。訂單表的重復(fù)問(wèn)題更有意思。當(dāng)時(shí)發(fā)現(xiàn)的重復(fù)訂單號(hào)第一次出現(xiàn)的金額是80元第二次是85元。這說(shuō)明上游系統(tǒng)對(duì)同一訂單做了多次修改并生成了新的快照記錄而不是真正的一單兩用。去重策略就不是簡(jiǎn)單保留任意一條而是要根據(jù)業(yè)務(wù)規(guī)則保留最后修改時(shí)間最晚的那條同時(shí)把兩條金額存進(jìn)一個(gè)數(shù)組字段方便后續(xù)排查。WITH deduped AS ( SELECT order_id, order_amount, create_time, ROW_NUMBER() OVER ( PARTITION BY order_id ORDER BY update_time DESC ) AS rn FROM ods_order ) SELECT * FROM deduped WHERE rn 1;這類去重最大的坑在于如果上游系統(tǒng)的更新邏輯本身有bug光靠數(shù)倉(cāng)去重是堵不住的。所以跑完清洗之后一定要回傳一個(gè)數(shù)據(jù)質(zhì)量異常報(bào)告給業(yè)務(wù)系統(tǒng)讓他們知道這里有重復(fù)產(chǎn)生的問(wèn)題。4.4 清洗規(guī)則發(fā)布與驗(yàn)證規(guī)則寫(xiě)完之后不能直接替換線上表。我習(xí)慣的發(fā)布流程是先做并行試跑把清洗后的數(shù)據(jù)落到一張新表dwd_order_clean_test和線上正在用的舊表做一次全量對(duì)比主鍵是否完全一致、關(guān)鍵指標(biāo)的差異量級(jí)是否在可接受范圍內(nèi)差異超過(guò)閾值比如訂單總額偏差超過(guò)5%就必須回頭查規(guī)則不能強(qiáng)行切換這一步非常關(guān)鍵。清洗規(guī)則本身有主觀判斷的成分這個(gè)字段置NULL還是填默認(rèn)值直接影響下游指標(biāo)。如果規(guī)則設(shè)計(jì)錯(cuò)了發(fā)布后報(bào)表異常半天時(shí)間就搭進(jìn)去了。那次項(xiàng)目的最終結(jié)果是訂單表的有效數(shù)據(jù)率從82%提升到了97%軌跡點(diǎn)漂移率從3.7%降到0.4%以下用戶次日留存指標(biāo)修正了將近1.2個(gè)百分點(diǎn)——這個(gè)修正幅度說(shuō)明之前的臟數(shù)據(jù)確實(shí)嚴(yán)重影響了業(yè)務(wù)判斷。5. 清洗后的質(zhì)量監(jiān)控規(guī)則老化、血緣追溯與回歸保障清洗規(guī)則上線不等于一勞永逸。數(shù)據(jù)是動(dòng)態(tài)的上游系統(tǒng)一改接口、產(chǎn)品一改邏輯你精心設(shè)計(jì)的清洗規(guī)則就可能失效。所以最后一塊內(nèi)容聊清洗后的質(zhì)量監(jiān)控問(wèn)題。5.1 監(jiān)控規(guī)則本身而不是監(jiān)控結(jié)果很多團(tuán)隊(duì)做數(shù)倉(cāng)質(zhì)量監(jiān)控光盯著今天的訂單量是否波動(dòng)超過(guò)10%這種做法太滯后。我習(xí)慣在DWD層直接埋規(guī)則校驗(yàn)點(diǎn)比如訂單表里status字段不在枚舉值范圍內(nèi)的記錄數(shù)必須為0手機(jī)號(hào)非法率不能超過(guò)0.1%訂單金額為負(fù)的記錄數(shù)必須為0當(dāng)日新增訂單的支付時(shí)間不能晚于次日0點(diǎn)這些校驗(yàn)點(diǎn)本質(zhì)上就是你清洗規(guī)則的鏡像。規(guī)則說(shuō)金額必須大于0那監(jiān)控就查金額小于等于0的記錄數(shù)。如果這個(gè)校驗(yàn)點(diǎn)爆了說(shuō)明上游出了問(wèn)題或者規(guī)則需要調(diào)整。5.2 血緣倒查規(guī)則變更影響面有多廣數(shù)倉(cāng)的血緣追溯在數(shù)據(jù)清洗這個(gè)語(yǔ)境下特別重要。當(dāng)你準(zhǔn)備改一條清洗規(guī)則時(shí)比如手機(jī)號(hào)非法時(shí)從置NULL改為填默認(rèn)值00000000000你必須立刻知道下游哪些表、哪些指標(biāo)會(huì)受影響。如果血緣關(guān)系不清晰你改了一條規(guī)則第二天一堆報(bào)表出問(wèn)題都不知道該找誰(shuí)。所以在設(shè)計(jì)數(shù)倉(cāng)時(shí)每個(gè)清洗任務(wù)我都要求必須有輸入表、輸出表、規(guī)則版本號(hào)三個(gè)元數(shù)據(jù)字段。出問(wèn)題的時(shí)候先查規(guī)則版本最近有沒(méi)有變過(guò)再倒查引用這張表的任務(wù)有哪些。這個(gè)習(xí)慣救過(guò)我很多次。5.3 定期回歸清洗邏輯也要做自動(dòng)化回歸最后是回歸保障。建議每隔一段時(shí)間我是按季度拉一批歷史數(shù)據(jù)用當(dāng)前版本的清洗規(guī)則重新跑一遍和上一版本的結(jié)果做對(duì)比。重點(diǎn)看兩件事有沒(méi)有規(guī)則改動(dòng)導(dǎo)致歷史數(shù)據(jù)結(jié)果漂移有沒(méi)有上游數(shù)據(jù)格式悄悄變化導(dǎo)致清洗命中率下降回歸通過(guò)之后清洗規(guī)則才算真正穩(wěn)定可靠。這個(gè)過(guò)程就像給數(shù)據(jù)質(zhì)量上了個(gè)保險(xiǎn)平時(shí)不起眼出了問(wèn)題才知道它的價(jià)值。做數(shù)據(jù)清洗這些年我最大的體會(huì)就是別把清洗當(dāng)成一個(gè)可以一步到位的動(dòng)作它和建模、監(jiān)控、血緣、元數(shù)據(jù)管理是糾纏在一起的。如果你只是按臨時(shí)需求東改一條西補(bǔ)一條那數(shù)據(jù)質(zhì)量永遠(yuǎn)在救火只有當(dāng)清洗規(guī)則變成數(shù)倉(cāng)建設(shè)的一等公民數(shù)據(jù)倉(cāng)庫(kù)才能真正成為值得信賴的分析底座。希望這篇文章能把你在數(shù)據(jù)清洗這件事上的思路捋順少踩幾個(gè)我踩過(guò)的坑。