倉SQL遷移Spark SQL:方言差異、UDF改造與數(shù)據(jù)驗(yàn)證實(shí)戰(zhàn)梳理)
從接手的那天起我就知道這活兒沒有表面看起來那么簡單。團(tuán)隊(duì)反饋過來的需求很簡短——把現(xiàn)有的Hive數(shù)倉SQL遷移到Spark SQL。我當(dāng)時(shí)的第一反應(yīng)是這不就是把SQL語句改一改、換個(gè)引擎跑嗎等真正干起來才發(fā)現(xiàn)整個(gè)spark-sql migration最耗時(shí)的地方根本不在翻譯SQL而在于兩套引擎在語義、元數(shù)據(jù)、UDF和運(yùn)行時(shí)行為上的隱性差異。這篇文章就把我在這個(gè)遷移項(xiàng)目里的完整經(jīng)歷、踩坑過程和沉淀下來的檢查清單一次說清楚。無論你是正準(zhǔn)備把Hive作業(yè)遷到Spark SQL還是想把老的Spark 2.x作業(yè)升級(jí)到Spark 3.x又或者是把其他引擎的SQL方言整理成Spark SQL規(guī)范都應(yīng)該能從這里面找到可以直接抄作業(yè)的部分。1. 動(dòng)手之前先分清你要做的是哪一種Spark SQL遷移很多人一聽到migration就默認(rèn)是在做引擎替換這其實(shí)是個(gè)誤區(qū)。我在項(xiàng)目啟動(dòng)前做的第一件事就是把遷移這件事拆成了三個(gè)完全不同的類型因?yàn)樗鼈兊墓ぷ髁俊L(fēng)險(xiǎn)點(diǎn)和驗(yàn)收方式完全不一樣。1.1 引擎替換型Hive作業(yè)遷到Spark SQL這是最典型的一種也就是把原本跑在Hive on MR或Hive on Tez上的SQL作業(yè)整體切換到Spark SQL引擎上執(zhí)行。這種遷移的核心矛盾在于Hive和Spark SQL雖然都支持絕大部分HiveQL語法但Spark SQL的Catalyst優(yōu)化器有自己的執(zhí)行語義尤其在子查詢、隱式類型轉(zhuǎn)換、NULL處理和窗口函數(shù)邊界上兩邊的行為差異非常明顯。我在項(xiàng)目里遇到的一個(gè)典型問題是Hive里允許把String類型字段直接和BigInt類型做等值JOIN引擎會(huì)自動(dòng)做寬松的隱式轉(zhuǎn)換跑起來一點(diǎn)問題沒有。但同樣的SQL切到Spark SQL如果不開兼容開關(guān)直接拋AnalysisException提示cannot resolve或類型不匹配。這類問題在存量幾百條SQL的數(shù)據(jù)倉庫里幾乎每十條就能碰上一次。1.2 版本升級(jí)型Spark 2.x作業(yè)遷到Spark 3.x第二種類型是Spark版本大版本升級(jí)。很多團(tuán)隊(duì)之前用Spark 2.4跑得好好的因?yàn)樾鹿δ?、安全補(bǔ)丁或云平臺(tái)服務(wù)期限的原因必須升到Spark 3.x。這類遷移表面上不動(dòng)SQL業(yè)務(wù)邏輯但Spark 3.0開始默認(rèn)啟用了ANSI模式的部分規(guī)則還把很多老參數(shù)標(biāo)記為deprecated導(dǎo)致原來能跑的生產(chǎn)作業(yè)突然在測(cè)試環(huán)境翻車。最典型的例子是spark.sql.legacy.allowCreatingManagedTableUsingNonemptyLocation這類遺留行為開關(guān)以及spark.sql.shuffle.partitions默認(rèn)值變化對(duì)性能的影響。這不是改SQL的問題而是調(diào)參和語義對(duì)齊的問題。1.3 方言轉(zhuǎn)換型把Presto、Flink SQL等方言整理為Spark SQL規(guī)范第三種類型相對(duì)少見但也不是沒有——把Presto/Trino或者Flink SQL的查詢邏輯統(tǒng)一改寫成Spark SQL規(guī)范。這種遷移的難點(diǎn)是函數(shù)名和語義映射比如Presto的approx_distinct到Spark SQL是approx_count_distinctFlink的HOP窗口在Spark SQL里要用window函數(shù)或flink風(fēng)格改寫稍不注意就會(huì)算出完全錯(cuò)誤的結(jié)果。MIGRATION COCKPIT操作手冊(cè)里有一句話我很認(rèn)同數(shù)據(jù)遷移的核心不是搬而是驗(yàn)證搬過去之后的行為一致。放在這里也一樣——搞清楚你是哪一種遷移才能決定后面每一步的精力分配。我在項(xiàng)目啟動(dòng)前先做了一輪資產(chǎn)盤點(diǎn)把待遷移SQL按來源、涉及引擎特性、是否含UDF做了分類這一步幫我在后面省下大量排查時(shí)間。2. 方言差異清單同一條SQL在兩個(gè)引擎里的行為可能完全不同如果說遷移項(xiàng)目的前半程是能跑通后半程就是跑得對(duì)。很多SQL在Hive里跑得很歡到Spark SQL里要么直接報(bào)錯(cuò)要么悄無聲息地給你一個(gè)不同的結(jié)果。這里我把自己實(shí)際遇到過且概率較高的方言差異點(diǎn)整理成了一份清單。2.1 ANSI模式開關(guān)一切差異的源頭Spark SQL從3.0開始引入spark.sql.ansi.enabled開啟后對(duì)類型轉(zhuǎn)換、除零、非法日期等行為的約束非常嚴(yán)格。Hive則默認(rèn)走盡力轉(zhuǎn)換的路線字符串轉(zhuǎn)數(shù)字轉(zhuǎn)不了就返回NULL而ANSI模式下直接拋異常。舉個(gè)例子在Hive里執(zhí)行SELECT CAST(abc AS INT)結(jié)果是NULL作業(yè)照常跑。同一個(gè)語句放到Spark SQL的ANSI模式下直接拋NumberFormatException。這會(huì)讓很多原本能用的臟數(shù)據(jù)邏輯突然變成作業(yè)失敗。我在遷移項(xiàng)目里的處理方案是對(duì)于存量Hive作業(yè)先用spark.sql.ansi.enabledfalse跑通再逐批開啟ANSI并修正SQL。這么做是為了把語法層遷移和數(shù)據(jù)質(zhì)量治理兩個(gè)目標(biāo)分開推進(jìn)避免混在一起后問題無法定位。2.2 常用函數(shù)行為差異對(duì)照表下面是我們?cè)谶w移過程中踩過且修復(fù)過的具體函數(shù)差異建議直接保存下來當(dāng)成排查手冊(cè)用場景Hive行為Spark SQL行為遷移建議字符串轉(zhuǎn)數(shù)字失敗返回NULL不報(bào)錯(cuò)ANSI開啟時(shí)報(bào)錯(cuò)關(guān)閉時(shí)返回NULL先判斷臟數(shù)據(jù)比例再用TRY_CAST做安全轉(zhuǎn)換datediff(end, start)返回INT天數(shù)返回INT天數(shù)但日期非法時(shí)行為不同統(tǒng)一先做日期合法性過濾substr與負(fù)索引負(fù)數(shù)索引從尾部計(jì)數(shù)負(fù)數(shù)索引返回空串改寫邏輯或統(tǒng)一約定參數(shù)非負(fù)count(distinct col)精確去重計(jì)數(shù)精確去重計(jì)數(shù)但大key傾斜風(fēng)險(xiǎn)更高數(shù)據(jù)量大時(shí)改用approx_count_distinctNULL排序升序時(shí)NULL默認(rèn)排在最前升序時(shí)NULL默認(rèn)排在最后可用NULLS FIRST/LAST控制顯式指定NULLS排序方向map與struct構(gòu)造語法較寬松對(duì)類型一致性要求更嚴(yán)格提前統(tǒng)一元素類型LATERAL VIEW explode空數(shù)組不輸出行空數(shù)組不輸出行但explode_outer才保留原行明確是否需要保留空行上面每一個(gè)差異都在我這次遷移中至少觸發(fā)過一次線上問題。特別是NULL排序這個(gè)點(diǎn)明明同一個(gè)結(jié)果集Hive和Spark排出來的前幾條記錄完全不同做數(shù)據(jù)驗(yàn)證時(shí)直接對(duì)不上賬。2.3 隱式類型轉(zhuǎn)換JOIN條件里的隱形殺手隱式類型轉(zhuǎn)換是遷移中最隱蔽的一類問題。在Hive里String和Int做等值JOIN時(shí)引擎會(huì)嘗試轉(zhuǎn)換通常在Hive里表現(xiàn)為把字符串先轉(zhuǎn)成數(shù)值。但Spark SQL在默認(rèn)情況下要求兩邊的數(shù)據(jù)類型匹配度高否則就走BroadcastNestedLoopJoin或者直接報(bào)錯(cuò)。我遇到過一個(gè)典型的生產(chǎn)故障兩張事實(shí)表用order_id關(guān)聯(lián)一張表的字段類型是INT另一張表在Hive建表時(shí)沒注意變成了STRING。Hive里跑了兩年都沒事遷移到Spark SQL后速度驟降排查發(fā)現(xiàn)是因?yàn)轭愋筒灰恢聦?dǎo)致JOIN沒法走SortMergeJoin而走了NestedLoop。修復(fù)方式也很簡單把STRING字段統(tǒng)一轉(zhuǎn)成INT或BIGINTJOIN性能馬上就恢復(fù)了。這里的經(jīng)驗(yàn)是遷移前先做一遍全量元數(shù)據(jù)對(duì)比把類型不一致的JOIN key單獨(dú)列出來不要等到作業(yè)上線了才通過性能問題反查。3. 表結(jié)構(gòu)與元數(shù)據(jù)遷移CREATE TABLE能跑通只是起點(diǎn)很多人在遷移時(shí)只盯著SQL語句忽略了一個(gè)更基礎(chǔ)的東西——目標(biāo)集群里的表定義是否和源集群完全一致。實(shí)際上表結(jié)構(gòu)層面的坑比SQL語法的坑更多而且一旦埋下后面所有作業(yè)都會(huì)受到影響。3.1 建表屬性檢查清單在Migrate your data類的操作手冊(cè)中元數(shù)據(jù)遷移永遠(yuǎn)排在第一步這里我把它落地成了一張可直接對(duì)照的表結(jié)構(gòu)差異檢查表存儲(chǔ)格式TextFile、SequenceFile、Parquet、ORC是否在目標(biāo)集群均已支持壓縮方式源表用Snappy、Zlib還是LZ4目標(biāo)集群的Codec是否配置齊全分區(qū)字段PARTITIONED BY的字段順序、類型是否完全一致分桶信息CLUSTERED BY和BUCKETS數(shù)量是否保留這直接影響Spark的bucket pruning表屬性tblproperties里的transient_lastDdlTime等參數(shù)是否需要保留字段注釋COMMENT信息是否遷移影響下游數(shù)據(jù)字典和數(shù)據(jù)治理行格式與SerDe自定義SerDe類在Spark SQL里是否可用我們項(xiàng)目的表都是從Hive Metastore遷移過來的理論上元數(shù)據(jù)應(yīng)該是同步的。但實(shí)際執(zhí)行時(shí)發(fā)現(xiàn)源集群里有一批用ROW FORMAT SERDE定義的表SerDe類只在老集群的Hive環(huán)境中注冊(cè)過Spark SQL執(zhí)行時(shí)無法加載導(dǎo)致一批作業(yè)全掛在表掃描階段。3.2 文件格式與壓縮Parquet和ORC不是隨便選的文件格式這塊我強(qiáng)烈建議在遷移前就統(tǒng)一好標(biāo)準(zhǔn)而不是沿用每個(gè)業(yè)務(wù)線各自的偏好。Parquet和ORC在Spark SQL里都支持得很好但要注意幾個(gè)細(xì)節(jié)Parquet的schema兼容性Spark SQL讀Parquet時(shí)如果文件里的字段類型和表定義不一致會(huì)走spark.sql.parquet.respectSummaryFiles和mergeSchema相關(guān)邏輯。如果老集群寫入了臟schema新集群可能需要開spark.sql.parquet.enableVectorizedReaderfalse才能正確讀取但這會(huì)犧牲性能。ORC的ACID支持如果源表是Hive ACID事務(wù)表尤其是帶UPDATE和DELETE的遷移到Spark SQL后要確認(rèn)版本。Spark 3.x支持讀取Hive ACID表但不支持所有的ACID寫入操作需要評(píng)估是否改寫寫入鏈路。壓縮算法的Codec名稱Hive里寫的是org.apache.hadoop.io.compress.SnappyCodecSpark SQL里經(jīng)常簡寫成snappy兩邊要能正確識(shí)別。如果Codec缺失作業(yè)不會(huì)報(bào)錯(cuò)但文件大小會(huì)異常膨脹。3.3 INSERT OVERWRITE的語義差異一個(gè)容易忽略的坑INSERT OVERWRITE在Hive和Spark SQL里的語義對(duì)分區(qū)表的處理不同。Hive的INSERT OVERWRITE默認(rèn)是覆蓋動(dòng)態(tài)分區(qū)對(duì)應(yīng)的分區(qū)目錄但Spark SQL在早期版本曾經(jīng)出現(xiàn)過覆蓋整張表或覆蓋非預(yù)期分區(qū)的行為。雖然Spark 3.x已經(jīng)默認(rèn)按分區(qū)覆蓋但如果你是從Spark 2.1或更早版本一路升級(jí)過來的一定要在測(cè)試環(huán)境驗(yàn)證一遍。我建議的統(tǒng)一規(guī)范是所有分區(qū)寫入都顯式指定靜態(tài)分區(qū)或者保證動(dòng)態(tài)分區(qū)的predicate足夠收斂不要依賴引擎默認(rèn)行為。這看起來保守但能在遷移期減少大量意外。4. UDF/UDAF遷移改造最容易拖慢進(jìn)度的隱蔽環(huán)節(jié)在Hive SQL向Spark SQL遷移的過程中最容易被低估的就是UDF和UDAF資產(chǎn)。很多團(tuán)隊(duì)的SQL看起來只是幾十個(gè)查詢但背地里掛了十幾個(gè)自定義函數(shù)有的還是用Hive的Java接口寫的。這些UDF在Spark SQL里能不能跑、性能如何直接決定整個(gè)遷移的排期。4.1 先盤點(diǎn)再動(dòng)手搞定UDF資產(chǎn)清單我接手項(xiàng)目時(shí)先做了一個(gè)UDF資產(chǎn)盤點(diǎn)按照以下維度建了一張表UDF名稱語言注冊(cè)方式邏輯復(fù)雜度能否用內(nèi)置函數(shù)替代遷移優(yōu)先級(jí)get_md5JavaHive永久函數(shù)低spark內(nèi)置md5可替代高parse_urlJavaHive永久函數(shù)低內(nèi)置parse_url可替代高json_extractPython臨時(shí)函數(shù)中建議改用get_json_object高geo_distanceJava永久UDF中無內(nèi)置替代中rank_by_groupScala永久UDAF高可用窗口函數(shù)改寫低這么一盤很多UDF其實(shí)是多余的。能用內(nèi)置函數(shù)替代的我建議一律用內(nèi)置函數(shù)因?yàn)閁DF在Spark SQL里不僅僅是寫法問題還涉及序列化、代碼生成和優(yōu)化器穿透性能差距很大。4.2 Hive UDF與Spark SQL的兼容調(diào)用Spark SQL天然支持加載Hive的UDF只需要在SQL里通過CREATE TEMPORARY FUNCTION或永久函數(shù)注冊(cè)即可。但要注意三點(diǎn)第一Jar包的依賴沖突。我們有一個(gè)Hive UDF依賴了舊版本的Guava和Jackson在Hive里沒事加載到Spark里直接和Spark自身依賴的版本沖突ClassNotFound和NoSuchMethodError滿天飛。這種問題排查周期很長強(qiáng)烈建議提前做依賴隔離測(cè)試。第二UDF的返回值類型聲明。有些Hive UDF用ObjectInspector實(shí)現(xiàn)返回值類型推斷在Spark SQL里可能會(huì)被識(shí)別為BinaryType而不是StringType導(dǎo)致結(jié)果集在展示和落盤時(shí)出現(xiàn)亂碼或者類型轉(zhuǎn)換異常。遇到這種情況需要在注冊(cè)函數(shù)時(shí)顯式指定返回類型。第三Python UDF的性能問題。Python UDF在Spark里每個(gè)分區(qū)的每行數(shù)據(jù)都要經(jīng)過Arrow序列化和反序列化數(shù)據(jù)量大時(shí)性能衰減非常明顯。我在項(xiàng)目里遇到一個(gè)Python UDF處理JSON字段源作業(yè)在Hive里跑15分鐘切到Spark后跑了將近兩個(gè)小時(shí)。后來把邏輯改寫為get_json_object和from_json內(nèi)置函數(shù)組合才把作業(yè)壓縮到10分鐘以內(nèi)。4.3 用SQL表達(dá)式替代UDF遷移之外的額外收益這里提供一個(gè)經(jīng)驗(yàn)思路遇到UDF先問自己三個(gè)問題——這個(gè)函數(shù)能不能用SQL表達(dá)式加內(nèi)置函數(shù)拼出來能不能用CASE WHEN解決能不能用窗口函數(shù)替代只有三個(gè)都回答不了才需要保留UDF。我用這種思路處理過一批自研字符串清洗函數(shù)它們?cè)驹贖ive里用Java實(shí)現(xiàn)邏輯是把一串文本里的手機(jī)號(hào)、郵箱、URL脫敏??雌饋砗軓?fù)雜但實(shí)際上用正則表達(dá)式regexp_replace配合幾個(gè)CASE WHEN就能實(shí)現(xiàn)不僅省去了UDF遷移的部署工作執(zhí)行速度還提升了4倍以上。5. 數(shù)據(jù)結(jié)果比對(duì)怎么確認(rèn)遷完之后的數(shù)是對(duì)的遷移到一半的時(shí)候業(yè)務(wù)方一定會(huì)問你一句你確定遷完之后的數(shù)據(jù)是對(duì)的嗎面對(duì)這種問題光拍胸脯沒有用得有一套數(shù)據(jù)驗(yàn)證的機(jī)制。我在這套機(jī)制上的投入幾乎和SQL改寫本身一樣多。5.1 三層校驗(yàn)法從粗到細(xì)逐步逼近我把數(shù)據(jù)比對(duì)設(shè)計(jì)成了三個(gè)層級(jí)從速度和置信度兩個(gè)維度做平衡第一層是行數(shù)與主鍵唯一性校驗(yàn)。最簡單也最可靠每個(gè)表跑一下COUNT(*)關(guān)鍵表再跑一下主鍵去重后的數(shù)量和源集群對(duì)比。如果連行數(shù)都對(duì)不上后面就不用比了。我們遷移的第一個(gè)月里有將近40%的表在這一層就被攔下來了原因包括分區(qū)目錄丟失、數(shù)據(jù)重復(fù)寫入和JOIN類型變化導(dǎo)致的行數(shù)膨脹。第二層是關(guān)鍵聚合指標(biāo)校驗(yàn)。對(duì)每個(gè)業(yè)務(wù)表抽取一組核心指標(biāo)比如金額求和、訂單數(shù)COUNT、用戶數(shù)COUNT(DISTINCT)在源集群和目標(biāo)集群分別執(zhí)行對(duì)比結(jié)果。這一層能發(fā)現(xiàn)大部分明細(xì)數(shù)據(jù)錯(cuò)亂的問題。需要注意COUNT(DISTINCT)在超大寬表上的執(zhí)行效率不高可以采樣后用approx_count_distinct對(duì)比精度完全夠用。第三層是抽樣明細(xì)比對(duì)。針對(duì)數(shù)據(jù)量特別大且不允許誤差的表我采用分層抽樣——按業(yè)務(wù)日期各抽一天的完整全量數(shù)據(jù)做全字段比對(duì)。比對(duì)方法是通過hash函數(shù)將每行序列化成指紋然后做集合差集。這一層的成本最高所以一般只應(yīng)用于核心交易表和用戶主數(shù)據(jù)表不會(huì)對(duì)全部表做。5.2 空值與邊界值數(shù)據(jù)比對(duì)中最容易忽略的部分?jǐn)?shù)據(jù)比對(duì)結(jié)果不一致很多時(shí)候不是真的一致性問題而是空值語義差異和邊界值格式化差異。舉個(gè)例子Hive里的空串和NULL是兩種不同的值但在下游分析時(shí)經(jīng)常被當(dāng)成等價(jià)處理。而Spark SQL在讀取某些字段時(shí)如果文件里實(shí)際是空串而表定義允許NULL部分寫入鏈路會(huì)把這個(gè)空串自動(dòng)轉(zhuǎn)成NULL。對(duì)比結(jié)果一跑差異值全來自這類字段。另外Decimal類型也有類似問題——源表和目標(biāo)表如果精度定義不一樣一個(gè)是DECIMAL(10,2)一個(gè)是DECIMAL(12,4)同一個(gè)數(shù)值在兩張表里Hash出來的指紋就不一樣。所以抽樣比對(duì)前必須先把兩端表的字段類型統(tǒng)一成同一套口徑。5.3 數(shù)據(jù)驗(yàn)證過程的自動(dòng)化手工執(zhí)行SQL驗(yàn)證在三五十張表的時(shí)候還能靠人肉推進(jìn)到了幾百張表的量級(jí)就完全不可行了。我建議把驗(yàn)證流程做成一個(gè)定時(shí)腳本用Shell或者Python按表分批執(zhí)行上述三層校驗(yàn)結(jié)果輸出到差異報(bào)告里。實(shí)操上有幾個(gè)細(xì)節(jié)可以供參考源集群和目標(biāo)集群的數(shù)據(jù)庫連接串、調(diào)度參數(shù)統(tǒng)一放在配置文件里避免腳本里硬編碼每張表的校驗(yàn)SQL由表結(jié)構(gòu)元數(shù)據(jù)自動(dòng)生成比如從Metastore讀取主鍵字段后自動(dòng)拼COUNT DISTINCT語句校驗(yàn)結(jié)果落庫每次遷移批次后自動(dòng)生成一個(gè)diff報(bào)告發(fā)送給相關(guān)方這套機(jī)制不復(fù)雜但能讓數(shù)據(jù)對(duì)不對(duì)從主觀判斷變成客觀可查的流程。6. 上線后踩過的坑三類高頻運(yùn)行錯(cuò)誤實(shí)錄遷移SQL全部跑通、數(shù)據(jù)比對(duì)全部通過之后并不代表項(xiàng)目結(jié)束了。作業(yè)到了生產(chǎn)環(huán)境在真實(shí)數(shù)據(jù)量和并發(fā)條件下還要面對(duì)運(yùn)行時(shí)的那一關(guān)。我把這個(gè)階段遇到的高頻錯(cuò)誤整理成三類每一類都附上了排查思路方便你將來遇到時(shí)能快速定位。6.1 運(yùn)行錯(cuò)誤一JOIN字段類型不一致導(dǎo)致的性能跳水現(xiàn)象是某個(gè)匯總作業(yè)從原來跑20分鐘變成了跑3小時(shí)以上而且集群資源占用居高不下。一開始我以為是資源競爭后來看Spark UI發(fā)現(xiàn)執(zhí)行計(jì)劃里出現(xiàn)了BroadcastNestedLoopJoin而正常的等值JOIN應(yīng)該走SortMergeJoin。排查時(shí)我打開了兩張表的schema才發(fā)現(xiàn)一張表的user_id是INT另一張表是STRING。在Hive里引擎自動(dòng)做了轉(zhuǎn)換所以察覺不出來到了Spark SQL由于類型不完全匹配優(yōu)化器選用了兜底的NestedLoop方案。修復(fù)方式很簡單把JOIN條件的字段顯式轉(zhuǎn)換成同一類型再查看執(zhí)行計(jì)劃確認(rèn)已經(jīng)恢復(fù)成SortMergeJoin。這件事也讓我養(yǎng)成了一個(gè)習(xí)慣任何涉及JOIN和分區(qū)的字段遷移前必須做一輪類型一致性檢查。6.2 運(yùn)行錯(cuò)誤二Decimal精度溢出現(xiàn)象是某個(gè)金額計(jì)算作業(yè)在遷移后偶發(fā)報(bào)錯(cuò)錯(cuò)誤信息是java.lang.ArithmeticException: Decimal overflow。源集群Hive里同樣的計(jì)算一直正常為什么到了Spark SQL就溢出原因是Spark SQL在計(jì)算DECIMAL(p,s)類型的乘法、除法時(shí)有一套嚴(yán)格的精度和小數(shù)位推導(dǎo)規(guī)則。比如兩個(gè)DECIMAL(10,2)相乘按規(guī)則會(huì)推導(dǎo)成DECIMAL(21,4)超出某些場景下的允許精度。Hive在對(duì)應(yīng)場景下則做的是更寬松的轉(zhuǎn)換。解決方案有兩個(gè)方向一是把相關(guān)字段的類型統(tǒng)一擴(kuò)展到更大精度二是修改計(jì)算邏輯中的CAST位置先轉(zhuǎn)成DOUBLE再計(jì)算并在最終結(jié)果處轉(zhuǎn)回DECIMAL。第二種方案要注意浮點(diǎn)數(shù)精度損失金額敏感場景建議先放大精度再運(yùn)算。6.3 運(yùn)行錯(cuò)誤三動(dòng)態(tài)分區(qū)暴增導(dǎo)致Driver OOM現(xiàn)象是某張分區(qū)表的寫入作業(yè)頻繁O(jiān)OM錯(cuò)誤日志定位到Driver端內(nèi)存不足。查看分區(qū)情況后發(fā)現(xiàn)這個(gè)作業(yè)寫入時(shí)按天和按渠道兩個(gè)字段做動(dòng)態(tài)分區(qū)某天數(shù)據(jù)量特別大生成了幾萬個(gè)分區(qū)Driver端維護(hù)分區(qū)元數(shù)據(jù)的內(nèi)存瞬間被打滿。這個(gè)問題的修復(fù)方案有三步第一步在寫入SQL里對(duì)不同數(shù)據(jù)量的渠道做拆分流量大的渠道走獨(dú)立作業(yè)單獨(dú)指定靜態(tài)分區(qū)寫入第二步動(dòng)態(tài)分區(qū)寫入前先預(yù)估分區(qū)數(shù)量超過閾值時(shí)改走批處理循環(huán)比如按渠道維度循環(huán)執(zhí)行一個(gè)個(gè)靜態(tài)分區(qū)寫入第三步適當(dāng)調(diào)大spark.sql.maxDynamicPartitions上限并增加Driver內(nèi)存。這里的核心經(jīng)驗(yàn)是**動(dòng)態(tài)分區(qū)不是越多越好它和Driver內(nèi)存之間有直接的線性關(guān)系。**在遷移大表時(shí)如果發(fā)現(xiàn)分區(qū)數(shù)異常增長應(yīng)該優(yōu)先從業(yè)務(wù)上拆分寫入任務(wù)而不是一味調(diào)大參數(shù)。7. 寫在最后遷移檢查清單與個(gè)人體會(huì)文章的最后一部分我按遷移前、遷移中、上線前、上線后四個(gè)階段整理了一份完整的檢查清單權(quán)當(dāng)項(xiàng)目復(fù)盤留檔。這份清單不一定適用于所有團(tuán)隊(duì)但大體框架可以復(fù)用。階段檢查項(xiàng)完成標(biāo)準(zhǔn)遷移前SQL資產(chǎn)盤點(diǎn)與UDF資產(chǎn)盤點(diǎn)輸出待遷SQL清單、UDF依賴清單遷移前元數(shù)據(jù)對(duì)比所有表/視圖字段類型、分區(qū)、存儲(chǔ)格式差異清單遷移前環(huán)境參數(shù)對(duì)齊ANSI模式、Codec、序列化器、Metastore版本梳理完畢遷移中SQL語法改造全部作業(yè)在測(cè)試集群跑通耗時(shí)記錄在案遷移中UDF替換與改寫無法替代的UDF已重新編譯并完成依賴隔離遷移中三層數(shù)據(jù)校驗(yàn)核心表行數(shù)、聚合、抽樣明細(xì)全部通過上線前執(zhí)行計(jì)劃對(duì)比JOIN類型、分區(qū)裁剪、Shuffle分區(qū)數(shù)符合預(yù)期上線前資源參數(shù)調(diào)優(yōu)shuffle分區(qū)數(shù)、Executor內(nèi)存、并行度按目標(biāo)數(shù)據(jù)量配置上線后灰度與監(jiān)控先切10%流量觀察再逐步放大到全量上線后回滾預(yù)案數(shù)據(jù)回流腳本和雙跑機(jī)制保留至少兩個(gè)賬期最后再分享一點(diǎn)個(gè)人體會(huì)。我剛開始做這個(gè)遷移項(xiàng)目的時(shí)候以為最大的困難會(huì)來自復(fù)雜的SQL改寫后來才發(fā)現(xiàn)真正消耗精力的是那些平時(shí)的隱形約定——隱式類型轉(zhuǎn)換、UDF的序列化行為、動(dòng)態(tài)分區(qū)數(shù)量這些細(xì)節(jié)平時(shí)在Hive里跑著一點(diǎn)感知都沒有換了引擎就全都冒出來了。如果你也在做類似的spark-sql migration我的建議是不要把遷移當(dāng)成一次性翻譯工程而是當(dāng)成一次數(shù)據(jù)鏈路的重構(gòu)。在排期里多預(yù)留兩到三成的時(shí)間給數(shù)據(jù)核對(duì)和異常排查這個(gè)時(shí)間會(huì)連本帶利地還給你。遷移本身不難難的是證明遷完之后行為和原來一致——在這一點(diǎn)上任何工具和腳本都替代不了你對(duì)自己數(shù)據(jù)的理解深度。