控系統(tǒng)架構(gòu)與落地實(shí)踐)
簡(jiǎn)介面向大數(shù)據(jù)金融信貸風(fēng)控領(lǐng)域?qū)W習(xí)者和畢業(yè)設(shè)計(jì)開發(fā)者的完整項(xiàng)目源碼包基于Hadoop與Spark技術(shù)棧實(shí)現(xiàn)信貸風(fēng)險(xiǎn)控制系統(tǒng)覆蓋數(shù)據(jù)接入、流式處理、風(fēng)控邏輯及可視化等環(huán)節(jié)適合課程設(shè)計(jì)、畢設(shè)或項(xiàng)目初期演示。壓縮包內(nèi)共69個(gè)文件以36個(gè)Java源碼、8個(gè)Scala源碼為主配合12個(gè)XML配置、5個(gè)properties屬性文件、SQL腳本及H5前端JS等整體約58KB工程結(jié)構(gòu)清晰便于按功能模塊查閱。已有325人學(xué)習(xí)下載代碼均測(cè)試運(yùn)行成功答辯平均分達(dá)96分并提供遠(yuǎn)程教學(xué)支持。參考README與數(shù)據(jù)庫(kù)腳本即可快速搭建環(huán)境重點(diǎn)理解Spark Streaming實(shí)時(shí)數(shù)據(jù)接入、MyBatis映射配置以及風(fēng)控流程設(shè)計(jì)等實(shí)踐要點(diǎn)下載后可直接導(dǎo)入IDE運(yùn)行調(diào)試便于二次開發(fā)。1. 基于Hadoop、Spark的信貸風(fēng)控系統(tǒng)在解決什么先看清它要扛下的數(shù)據(jù)量做信貸業(yè)務(wù)的同學(xué)都熟悉這個(gè)畫面一天涌入幾十萬(wàn)筆申請(qǐng)每筆申請(qǐng)背后是用戶基本信息、多頭借貸記錄、設(shè)備指紋、APP行為軌跡、第三方黑名單光原始字段就上百個(gè)。傳統(tǒng)MySQL單庫(kù)撐死幾千QPS十幾張表做一次join直接超時(shí)更別提把一年甚至三年的歷史數(shù)據(jù)拉出來(lái)做特征回溯?;贖adoop、Spark的大數(shù)據(jù)金融信貸風(fēng)控系統(tǒng)核心就是把數(shù)據(jù)存儲(chǔ)和特征計(jì)算搬到分布式框架上HDFS承接海量明細(xì)數(shù)據(jù)Spark跑離線特征工程、批量評(píng)分和模型訓(xùn)練再配合規(guī)則引擎處理實(shí)時(shí)審批決策。它解決的是三個(gè)具體問題——數(shù)據(jù)存不下、特征算不動(dòng)、風(fēng)險(xiǎn)評(píng)分出不來(lái)。適合正在做信貸技術(shù)選型的工程師也適合高校大數(shù)據(jù)方向拿一個(gè)完整項(xiàng)目做課程設(shè)計(jì)的同學(xué)。你先搞清楚這套系統(tǒng)的數(shù)據(jù)鏈路和關(guān)鍵取舍再去看源碼和文檔效率會(huì)高很多。2. 架構(gòu)設(shè)計(jì)與技術(shù)選型為什么是Hadoop加Spark而不是其他組合2.1 先做減法哪些方案被排除理由是什么常見的第一反應(yīng)是用MySQL加定時(shí)任務(wù)搞定一切。但在信貸風(fēng)控場(chǎng)景里原始數(shù)據(jù)里有用戶授權(quán)讀取的運(yùn)營(yíng)商通話詳單、電商消費(fèi)記錄、社保公積金流水單用戶每月產(chǎn)生的明細(xì)記錄就有幾百條百萬(wàn)用戶就是幾億行。MySQL單表過億之后即使加了索引復(fù)雜的聚合查詢也要幾十秒甚至分鐘級(jí)根本撐不住風(fēng)控調(diào)參時(shí)反反復(fù)復(fù)的特征回溯。純Flink流處理方案也是一個(gè)選擇但風(fēng)控場(chǎng)景里真正耗費(fèi)算力的不是實(shí)時(shí)計(jì)算本身而是全量數(shù)據(jù)的批量特征計(jì)算、模型訓(xùn)練和離線回測(cè)。Flink擅長(zhǎng)的是秒級(jí)窗口內(nèi)的實(shí)時(shí)計(jì)算把這部分任務(wù)硬塞給Flink集群成本會(huì)翻很多倍而且Flink對(duì)狀態(tài)后端的管理也比Spark的批處理模型復(fù)雜。Hadoop加Spark是更成熟的配合方式HDFS做分布式存儲(chǔ)底座Spark跑內(nèi)存計(jì)算MapReduce留作極重的歷史回溯兜底任務(wù)。這套組合在金融行業(yè)被驗(yàn)證了超過十年招聘市場(chǎng)上也確實(shí)大量崗位要求這兩項(xiàng)技術(shù)。2.2 總體架構(gòu)五層數(shù)據(jù)流和一個(gè)核心整套系統(tǒng)的數(shù)據(jù)流可以拆成五個(gè)層級(jí)采集層Flume采集服務(wù)器日志Kafka承接業(yè)務(wù)系統(tǒng)實(shí)時(shí)推送的申請(qǐng)事件兩份數(shù)據(jù)最終都落到HDFS。存儲(chǔ)層HDFS存放原始日志和數(shù)倉(cāng)各層明細(xì)數(shù)據(jù)HBase用來(lái)存實(shí)時(shí)特征結(jié)果和反欺詐名單支持毫秒級(jí)查詢。計(jì)算層Spark負(fù)責(zé)離線批處理、特征工程、模型訓(xùn)練Hadoop MapReduce處理極端耗時(shí)的全量回溯任務(wù)。服務(wù)層規(guī)則引擎加載黑名單、硬性規(guī)則進(jìn)行第一輪攔截評(píng)分模型對(duì)通過規(guī)則的申請(qǐng)輸出信用評(píng)分封裝成REST接口。展示層運(yùn)營(yíng)后臺(tái)和大屏展示每日申請(qǐng)量、通過率、逾期率、評(píng)分分布等核心指標(biāo)。整個(gè)架構(gòu)的核心是中間那層特征寬表和評(píng)分模型。沒有特征寬表Spark的算力無(wú)處安放沒有評(píng)分模型前面的存儲(chǔ)和計(jì)算只是存了一堆用不上的數(shù)據(jù)。這也是后面兩章要重點(diǎn)展開的部分。2.3 集群部署起點(diǎn)偽分布式搭建、HA 架構(gòu)和生產(chǎn)集群規(guī)劃學(xué)習(xí)階段不建議一上來(lái)就搞多節(jié)點(diǎn)集群。先在單機(jī)做Hadoop偽分布式搭建把HDFS的NameNode和DataNode、YARN的ResourceManager和NodeManager之間的關(guān)系跑明白再用Spark的local模式提交幾個(gè)作業(yè)理解存儲(chǔ)和計(jì)算的協(xié)作邏輯。集群里還有一個(gè)關(guān)鍵組件是ZookeeperHadoop和Zookeeper整合實(shí)戰(zhàn)解決的核心問題是HA模式下NameNode的自動(dòng)故障切換——主NameNode掛了備用節(jié)點(diǎn)要能自動(dòng)頂上否則整個(gè)HDFS就癱了。生產(chǎn)環(huán)境必須上HA架構(gòu)這一點(diǎn)沒有任何商量余地。NameNode是HDFS的單點(diǎn)元數(shù)據(jù)都在它內(nèi)存里一掛全掛。Zookeeper負(fù)責(zé)協(xié)調(diào)兩個(gè)NameNode的主備狀態(tài)JournalNode負(fù)責(zé)同步編輯日志。生產(chǎn)集群的規(guī)模按數(shù)據(jù)量倒推1000萬(wàn)級(jí)注冊(cè)用戶每天新增日志約200GBKeep一個(gè)50個(gè)節(jié)點(diǎn)左右的集群存儲(chǔ)和計(jì)算基本夠用。節(jié)點(diǎn)規(guī)格建議是每臺(tái)32GB內(nèi)存、8核CPU、4塊4TB硬盤這樣的配置在大多數(shù)信貸業(yè)務(wù)里能從業(yè)務(wù)初期撐到中期。規(guī)模再大的話要考慮的是Spark任務(wù)的資源隔離而不是繼續(xù)無(wú)限加節(jié)點(diǎn)。3. 數(shù)據(jù)接入與存儲(chǔ)從原始JSON到可查詢的特征寬表3.1 Spark中讀取JSON日志的兩種姿勢(shì)和一個(gè)關(guān)鍵坑業(yè)務(wù)系統(tǒng)上報(bào)的日志絕大多數(shù)是JSON格式每條申請(qǐng)記錄一個(gè)JSON文件或者一行一個(gè)JSON。Spark讀取JSON最常見的錯(cuò)誤是讓Spark自己推斷schema數(shù)據(jù)量大時(shí)這個(gè)推斷過程會(huì)觸發(fā)額外的掃描而且碰到包含多種字段形態(tài)的數(shù)據(jù)時(shí)推斷結(jié)果經(jīng)常和你預(yù)期不符——比如金額字段有時(shí)是字符串有時(shí)是數(shù)字Spark默認(rèn)推斷成string或bigint取出來(lái)才發(fā)現(xiàn)類型不對(duì)。我一般會(huì)直接在讀取時(shí)顯式指定schema代價(jià)是維護(hù)一個(gè)schema定義但換來(lái)的是穩(wěn)定性和可預(yù)測(cè)性。典型的讀取代碼長(zhǎng)這樣from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, TimestampType spark SparkSession.builder \ .appName(credit_json_etl) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate() # 顯式定義schema避免Spark自己推導(dǎo)導(dǎo)致的類型漂移 schema StructType([ StructField(user_id, StringType(), True), StructField(apply_time, TimestampType(), True), StructField(loan_amount, DoubleType(), True), StructField(device_id, StringType(), True), StructField(channel, StringType(), True), StructField(extra_info, StringType(), True) # 原始字段先保留后續(xù)解析 ]) # 讀取HDFS上的原始JSON目錄按天分區(qū) raw_df spark.read \ .option(multiline, false) \ .schema(schema) \ .json(/data/raw/credit_apply/dt2024-01-01/) raw_df.printSchema() raw_df.show(5, truncateFalse)這段代碼里multiline參數(shù)要按實(shí)際數(shù)據(jù)格式設(shè)置。每行一個(gè)JSON對(duì)象時(shí)設(shè)為false整個(gè)文件是一個(gè)大JSON數(shù)組時(shí)設(shè)為true。shuffle.partitions設(shè)置成200是一個(gè)比較保守的起步值它控制了Spark執(zhí)行shuffle操作時(shí)的分區(qū)數(shù)量數(shù)據(jù)量大的時(shí)候這個(gè)值設(shè)置得太小單個(gè)任務(wù)處理的數(shù)據(jù)過多很容易撐爆executor內(nèi)存。3.2 數(shù)倉(cāng)分層ODS、DWD、ADS三層的建模思路原始數(shù)據(jù)直接拿來(lái)算特征是不現(xiàn)實(shí)的JSON里嵌套的字段解析要消耗大量計(jì)算資源而且重復(fù)解析會(huì)有一致性問題。標(biāo)準(zhǔn)做法是建立三層數(shù)倉(cāng)ODS層原樣保留原始JSONDWD層做清洗、脫敏、解析字段ADS層做特征寬表。ODS層就是上面讀取的原始數(shù)據(jù)不做任何加工只按天分區(qū)。DWD層用Spark SQL做清洗和解析典型做法是用get_json_object把JSON里的嵌套字段提取出來(lái)同時(shí)對(duì)身份證號(hào)、手機(jī)號(hào)做脫敏處理。這一步很關(guān)鍵信貸數(shù)據(jù)涉及個(gè)人敏感信息審計(jì)時(shí)會(huì)要求你證明原始數(shù)據(jù)沒有直接暴露在計(jì)算鏈路里。DWD層建表語(yǔ)句參考CREATE EXTERNAL TABLE IF NOT EXISTS dwd_credit_apply ( user_id STRING COMMENT 用戶ID, apply_time TIMESTAMP COMMENT 申請(qǐng)時(shí)間, loan_amount DOUBLE COMMENT 申請(qǐng)金額, loan_term INT COMMENT 申請(qǐng)期限月, device_brand STRING COMMENT 設(shè)備品牌, channel_code STRING COMMENT 渠道編碼, city_level INT COMMENT 城市等級(jí) 1-5, -- 脫敏后的手機(jī)號(hào)只保留前3后4 mobile_masked STRING COMMENT 脫敏手機(jī)號(hào) ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION /warehouse/dwd/credit_apply; -- 按天覆蓋寫入分區(qū) INSERT OVERWRITE TABLE dwd_credit_apply PARTITION (dt2024-01-01) SELECT user_id, apply_time, CAST(loan_amount AS DOUBLE), CAST(loan_term AS INT), get_json_object(extra_info, $.device_brand) AS device_brand, get_json_object(extra_info, $.channel_code) AS channel_code, get_json_object(extra_info, $.city_level) AS city_level, concat(substr(mobile, 1, 3), ****, substr(mobile, 8, 4)) AS mobile_masked FROM /data/raw/credit_apply/dt2024-01-01;建表時(shí)用STORED AS PARQUET的列式存儲(chǔ)查詢時(shí)只讀取需要的列IO量可以降一半以上對(duì)特征計(jì)算階段的多列讀取特別友好。INSERT OVERWRITE的寫法保證分區(qū)內(nèi)數(shù)據(jù)是冪等的重跑任務(wù)不會(huì)產(chǎn)生重復(fù)數(shù)據(jù)。3.3 特征寬表為什么要寬以及怎么join才不翻車特征計(jì)算階段最忌諱的是每次評(píng)分都去臨時(shí)join十幾張表。業(yè)界標(biāo)準(zhǔn)做法是預(yù)先算好一張?zhí)卣鲗挶怼恳恍惺且粋€(gè)用戶每一列是一個(gè)特征。寬表的好處在于訓(xùn)練模型時(shí)直接把這張表喂給機(jī)器學(xué)習(xí)算法不需要再處理join邏輯實(shí)時(shí)評(píng)分時(shí)查詢一次就能拿到一個(gè)用戶的全部特征。構(gòu)建寬表最常見的坑是join數(shù)據(jù)傾斜。比如用用戶維度join授信記錄某幾個(gè)用戶借款次數(shù)上千次會(huì)導(dǎo)致對(duì)應(yīng)reduce任務(wù)數(shù)據(jù)量遠(yuǎn)大于其他任務(wù)表現(xiàn)為整體任務(wù)卡在99%。緩解手段有三個(gè)先把用戶維度的數(shù)據(jù)做聚合壓到每用戶一條對(duì)join key做加鹽處理分散熱點(diǎn)或者干脆把低頻維度做成廣播變量。第三個(gè)手段在Spark里操作最簡(jiǎn)單核心代碼如下from pyspark.sql import functions as F # 用戶基本信息表數(shù)據(jù)量在幾十萬(wàn)量級(jí)可以廣播 user_info spark.table(dwd_user_info) user_features spark.table(dwd_credit_apply).groupBy(user_id).agg( F.count(user_id).alias(apply_cnt), F.avg(loan_amount).alias(avg_loan_amount), F.max(loan_amount).alias(max_loan_amount) ) # 廣播小表避免shuffle user_features_with_info user_features.join( broadcast(user_info), onuser_id, howleft ) # 覆蓋寫入ADS特征寬表按天全量刷新 user_features_with_info.write \ .mode(overwrite) \ .format(parquet) \ .save(/warehouse/ads/user_feature_wide_table)這段代碼里broadcast會(huì)把小表分發(fā)到每個(gè)executor的內(nèi)存中join時(shí)不需要做shuffle這是Spark里最廉價(jià)的性能優(yōu)化手段之一。寬表的刷新策略按業(yè)務(wù)容忍度來(lái)定信貸場(chǎng)景通常T1刷新就行——當(dāng)天凌晨用前一天的數(shù)據(jù)重算全量寬表第二天白天做實(shí)時(shí)查詢時(shí)讀的就是最新版本。4. 特征工程與風(fēng)險(xiǎn)模型把Spark算力變成授信分4.1 特征計(jì)算與清洗Spark DataFrame的標(biāo)準(zhǔn)化數(shù)據(jù)處理特征寬表建好之后下一步是特征工程——把原始字段變成模型能用的數(shù)值特征。信貸風(fēng)控常見的特征類型包括用戶基本屬性年齡、城市等級(jí)、職業(yè)類別、借款行為申請(qǐng)次數(shù)、平均借款金額、借貸間隔、設(shè)備信息設(shè)備使用時(shí)長(zhǎng)、越獄/root標(biāo)記、外部征信數(shù)據(jù)逾期次數(shù)、查詢次數(shù)。這些特征不能直接進(jìn)模型要做三類處理缺失值填充、異常值截?cái)?、?shù)值標(biāo)準(zhǔn)化。Spark DataFrame做這套處理非常順手代碼示例如下from pyspark.sql import functions as F from pyspark.sql.functions import col # 讀取特征寬表 feature_df spark.read.parquet(/warehouse/ads/user_feature_wide_table) # 填充缺失值數(shù)值列用中位數(shù)類別列用unknown stat feature_df.select( F.expr(percentile_approx(age, 0.5)).alias(age_median) ).collect()[0] feature_df feature_df.fillna({ age: stat[age_median], city_level: 3, # 缺失城市等級(jí)默認(rèn)按三線城市處理 device_use_days: 0, channel_code: unknown }) # 異常值截?cái)喑^99分位數(shù)的值直接截?cái)喾乐箻O端值拉偏模型 quantile_age feature_df.approxQuantile(age, [0.99], 0.01)[0] quantile_amount feature_df.approxQuantile(loan_amount, [0.99], 0.01)[0] feature_df feature_df.withColumn( age, F.least(col(age), F.lit(quantile_age)) ).withColumn( loan_amount, F.least(col(loan_amount), F.lit(quantile_amount)) ) # 標(biāo)準(zhǔn)化Z-score讓不同量綱的特征處于同一尺度 mu_age, std_age feature_df.select( F.mean(age).alias(mu), F.stddev(age).alias(std) ).collect()[0] feature_df feature_df.withColumn( age_zscore, (col(age) - mu_age) / std_age )approxQuantile是Spark提供的近似分位數(shù)計(jì)算它不需要全量排序用抽樣估算對(duì)于分位數(shù)這種統(tǒng)計(jì)量精度完全夠。標(biāo)準(zhǔn)化用的是Z-score對(duì)于邏輯回歸這類線性模型是必須的不標(biāo)準(zhǔn)化的話數(shù)值大的特征會(huì)在梯度計(jì)算中占主導(dǎo)地位。對(duì)于樹模型標(biāo)準(zhǔn)化非必需但做了也沒有副作用所以我一般統(tǒng)一做掉。4.2 規(guī)則引擎和評(píng)分卡模型怎么配合有了特征接下來(lái)是關(guān)鍵問題規(guī)則引擎和評(píng)分卡模型怎么分工。規(guī)則引擎處理那些一刀切的硬性風(fēng)險(xiǎn)——命中黑名單直接拒、申請(qǐng)金額超過限額直接拒、同一設(shè)備一天申請(qǐng)超過5次直接拒。規(guī)則引擎速度快、可解釋性強(qiáng)、修改即時(shí)生效適合做第一道攔截。評(píng)分卡模型則處理那些介于中間的申請(qǐng)用一個(gè)分?jǐn)?shù)來(lái)度量違約概率。評(píng)分卡最底層的模型通常是邏輯回歸因?yàn)榫€性模型天然具備可解釋性。在實(shí)際落地中為了更好處理非線性關(guān)系會(huì)先對(duì)連續(xù)特征做WOE分箱。這里用pyspark.ml.feature里的相關(guān)組件來(lái)實(shí)現(xiàn)from pyspark.ml.feature import QuantileDiscretizer # 對(duì)連續(xù)特征做分箱讓模型學(xué)習(xí)非線性關(guān)系 discretizer QuantileDiscretizer( numBuckets10, inputColage, outputColage_bucket ) # 分箱結(jié)果轉(zhuǎn)成OneHot編碼后輸入邏輯回歸 from pyspark.ml.feature import OneHotEncoder ohe OneHotEncoder( inputCols[age_bucket, city_level_bucket], outputCols[age_ohe, city_ohe] ) # 組裝特征向量并訓(xùn)練邏輯回歸 from pyspark.ml.classification import LogisticRegression from pyspark.ml.feature import VectorAssembler assembler VectorAssembler( inputCols[age_ohe, city_ohe, apply_cnt, avg_loan_amount], outputColfeatures_vector ) lr LogisticRegression( featuresColfeatures_vector, labelColis_default, maxIter100, regParam0.01 )訓(xùn)練做完之后把模型的系數(shù)換算成分?jǐn)?shù)每個(gè)分箱對(duì)應(yīng)一個(gè)分?jǐn)?shù)加總后映射到300-900分的信用分區(qū)間。分?jǐn)?shù)越高違約概率越低。這個(gè)分?jǐn)?shù)段業(yè)內(nèi)默認(rèn)是300分以下堅(jiān)決拒、650分以上直接放、中間走人工復(fù)核。模型訓(xùn)練迭代的次數(shù)不用追求太多信貸場(chǎng)景里模型更新周期通常是季度級(jí)別特征是主導(dǎo)模型算法排第二。4.3 讓評(píng)分接口算得快Spark內(nèi)存參數(shù)與提交配置評(píng)分模型落地到線上服務(wù)常見的方式是把模型跑在Spark Streaming或者Structured Streaming上接Kafka里的申請(qǐng)事件流式計(jì)算實(shí)時(shí)給分。這時(shí)候性能參數(shù)的重要性不亞于模型本身的準(zhǔn)確率。Spark Streaming消費(fèi)Kafka數(shù)據(jù)做評(píng)分的典型配置參數(shù)如下from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(credit_scoring_streaming) \ .config(spark.executor.memory, 4g) \ .config(spark.executor.cores, 4) \ .config(spark.sql.shuffle.partitions, 400) \ .config(spark.streaming.kafka.maxRatePerPartition, 1000) \ .config(spark.sql.streaming.checkpointLocation, /warehouse/checkpoint/credit_scoring) \ .getOrCreate()spark.executor.memory決定了每個(gè)executor的JVM堆內(nèi)存4g是生產(chǎn)環(huán)境的保守起步值別忘了堆外內(nèi)存實(shí)際申請(qǐng)的物理內(nèi)存要比這個(gè)值多20%左右。maxRatePerPartition限制每個(gè)Kafka分區(qū)每秒最多消費(fèi)1000條防止流量突增直接把集群打爆。checkpointLocation一定要配流式任務(wù)的進(jìn)度元數(shù)據(jù)全部存在這里不配的話任務(wù)重啟后會(huì)丟數(shù)據(jù)或者重復(fù)消費(fèi)。提交作業(yè)的方式也要配套。推薦用spark-submit提交到Y(jié)ARN集群而不是直接spark-submit跑在本地模式。提交命令參考spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 4 \ --num-executors 30 \ --conf spark.sql.shuffle.partitions400 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0 \ credit_scoring_streaming.pynum-executors乘上之前配置的核數(shù)控制著整個(gè)作業(yè)的并行度。核數(shù)從4往上加收益會(huì)快速遞減因?yàn)閠ask調(diào)度和網(wǎng)絡(luò)通信的開銷也在增長(zhǎng)。對(duì)大部分信貸數(shù)據(jù)量級(jí)來(lái)說30個(gè)executor已經(jīng)能滿足秒級(jí)評(píng)分的要求。如果評(píng)分延遲還是壓不下來(lái)先看特征寬表有沒有觸發(fā)不必要的shuffle再考慮加資源這兩者的順序千萬(wàn)別搞反。5. 避坑清單Hadoop和Spark集群落地最常見的5個(gè)翻車現(xiàn)場(chǎng)5.1 數(shù)據(jù)傾斜導(dǎo)致某個(gè)executor內(nèi)存溢出現(xiàn)象Spark作業(yè)卡在某個(gè)stage 99%不動(dòng)日志刷出Container killed on request. Exit code is 143或者直接拋OOM。界面上看有一部分task執(zhí)行時(shí)間異常長(zhǎng)大部分task早就跑完了。原因特征寬表或者中間結(jié)果里某個(gè)key的數(shù)據(jù)量遠(yuǎn)大于其他key。信貸數(shù)據(jù)里最常見的傾斜源是user_id——部分高頻用戶或者渠道商賬號(hào)對(duì)應(yīng)的數(shù)據(jù)量是普通用戶的幾百倍這些數(shù)據(jù)全部落到同一個(gè)分區(qū)把對(duì)應(yīng)的executor打爆。解決優(yōu)先對(duì)傾斜key做過濾或單獨(dú)處理比如單獨(dú)處理借款次數(shù)超過100的用戶如果傾斜程度可控給Spark配置加spark.sql.adaptive.enabledtrue和spark.sql.adaptive.skewJoin.enabledtrueSpark3.0以上的版本開啟后會(huì)自動(dòng)拆分傾斜分區(qū)最后的手段才是加內(nèi)存。這一步的排查速度決定了你的加班時(shí)長(zhǎng)我一般會(huì)把傾斜前后的數(shù)據(jù)量打出來(lái)對(duì)比五分鐘定位。5.2 讓Spark自己推斷JSON的schema慢到懷疑人生現(xiàn)象讀取一個(gè)幾十GB的JSON目錄Spark作業(yè)在讀取階段卡了很久日志顯示一直在掃描文件。原因Spark為了推斷出所有字段的類型需要先完整掃描一遍數(shù)據(jù)然后再真正讀取一遍總共兩遍IO。目錄里文件數(shù)量多、嵌套層級(jí)深的時(shí)候時(shí)間成本非常可觀。另一個(gè)副作用是推斷結(jié)果不穩(wěn)定數(shù)據(jù)形態(tài)稍微一變字段類型就跟著漂。解決第3章已經(jīng)寫過讀取JSON時(shí)用.schema()顯式指定。維護(hù)schema確實(shí)麻煩但比每次任務(wù)跑兩遍要值得多。如果JSON里字段特別多可以先跑一次小樣本推斷出schema然后微調(diào)類型固化到代碼里當(dāng)常量。5.3 小文件過多把NameNode壓垮現(xiàn)象集群沒跑什么大任務(wù)但HDFS的NameNode頻繁告警RPC響應(yīng)延遲飆到幾秒甚至幾十秒。原因上游采集任務(wù)每個(gè)批次只寫幾十MB一天下來(lái)在HDFS上生成了上萬(wàn)個(gè)小文件。NameNode把每個(gè)文件、每個(gè)block的元數(shù)據(jù)都放在內(nèi)存里文件數(shù)量一多內(nèi)存占用上去讀寫請(qǐng)求的響應(yīng)性能直線下降。Spark寫特征寬表時(shí)如果分區(qū)設(shè)置過多且數(shù)據(jù)量不大也會(huì)產(chǎn)生同樣的問題。解決一是上游把小塊數(shù)據(jù)合并成大文件再上傳控制單文件在128MB以上二是對(duì)已有的小文件跑一次合并任務(wù)。Spark側(cè)的設(shè)置是控制輸出分區(qū)數(shù)核心代碼是這樣# 寫特征寬表前先重分區(qū)控制輸出文件數(shù)量 from pyspark.sql import functions as F feature_df \ .repartition(50) \ .write \ .mode(overwrite) \ .format(parquet) \ .save(/warehouse/ads/user_feature_wide_table)repartition(50)把輸出分片收縮到50個(gè)左右每個(gè)文件大約128MB這是一個(gè)平衡NameNode壓力和后續(xù)讀取并行度的折中值。文件太大讀取并行度會(huì)下降太小元數(shù)據(jù)壓力又上來(lái)。5.4 Zookeeper超時(shí)導(dǎo)致HA主備反復(fù)切換現(xiàn)象集群狀態(tài)不穩(wěn)定NameNode頻繁自動(dòng)切換甚至出現(xiàn)雙主或同時(shí)無(wú)人服務(wù)的情況業(yè)務(wù)端隔幾分鐘就報(bào)一次HDFS連接失敗。原因Hadoop和Zookeeper整合實(shí)戰(zhàn)里最常見的一個(gè)坑——NameNode和Zookeeper之間的心跳超時(shí)設(shè)置過短。Zookeeper需要同時(shí)維護(hù)NameNode元數(shù)據(jù)、HBase協(xié)調(diào)、Kafka協(xié)調(diào)多個(gè)角色的會(huì)話任何一次GC暫?;蛘呔W(wǎng)絡(luò)抖動(dòng)超過超時(shí)閾值Zookeeper就會(huì)判定NameNode失聯(lián)觸發(fā)自動(dòng)切換。切換本身又是高開銷操作元數(shù)據(jù)加載需要時(shí)間于是形成切換-搶主-再切換的惡性循環(huán)。解決把zookeeper.session.timeout從默認(rèn)的10秒左右調(diào)大到30秒dfs.namenode.avoid.read.stale.datanode加上讓NameNode對(duì)DataNode的陳舊狀態(tài)容忍度高一些。調(diào)參后的自檢方法是殺掉active節(jié)點(diǎn)的進(jìn)程觀察standby節(jié)點(diǎn)能否在1分鐘內(nèi)接管然后恢復(fù)服務(wù)。這個(gè)動(dòng)作要在業(yè)務(wù)低峰期做不然一個(gè)誤判就直接影響線上審批鏈路。5.5 物理內(nèi)存和容器內(nèi)存對(duì)不上任務(wù)被莫名殺掉現(xiàn)象Spark任務(wù)運(yùn)行到一半executor被YARN強(qiáng)制殺掉錯(cuò)誤日志只有一行Container killed by YARN for exceeding memory limits。明明executor-memory配的不高怎么還被殺。原因Spark申請(qǐng)的內(nèi)存默認(rèn)只算JVM堆內(nèi)內(nèi)存但實(shí)際運(yùn)行時(shí)還有堆外內(nèi)存、線程棧、網(wǎng)絡(luò)緩沖、Python進(jìn)程如果用PySpark占用的內(nèi)存。YARN限制的是整個(gè)容器進(jìn)程的物理內(nèi)存上限。申請(qǐng)4g的時(shí)候?qū)嶋H使用可能已經(jīng)超過5g超過容器上限直接被kill。解決申請(qǐng)內(nèi)存時(shí)預(yù)留20%-30%的buffer給堆外和Python進(jìn)程比如需要4g的堆內(nèi)就配spark.executor.memory4g的同時(shí)配spark.executor.memoryOverhead2g。YARN側(cè)還要開啟內(nèi)存檢測(cè)的寬松模式設(shè)置yarn.nodemanager.vmem-check-enabledfalse只檢查物理內(nèi)存不查虛擬內(nèi)存。這套配置在很多博客里都語(yǔ)焉不詳屬于那種你以為配好了實(shí)際上沒配的玄學(xué)參數(shù)。6. 拿到源代碼和文檔后怎么用先跑通最小閉環(huán)再改出你自己的風(fēng)控系統(tǒng)拿到這套基于Hadoop、Spark的信貸風(fēng)控系統(tǒng)源代碼我建議你按三條線往下推不要上來(lái)就翻模型代碼。第一條線是跑通數(shù)據(jù)鏈路——把生產(chǎn)環(huán)境的數(shù)據(jù)源換成你自己的樣本數(shù)據(jù)哪怕是模擬的100萬(wàn)條先跟著文檔把ODS到DWD的清洗任務(wù)跑完確認(rèn)HDFS上能看到正確的分區(qū)數(shù)據(jù)。第二條線是跑通Spark作業(yè)——把特征寬表的生成腳本執(zhí)行一遍看輸出和文檔里截圖的數(shù)量級(jí)是否對(duì)得上這一步最好在偽分布式或單機(jī)測(cè)試集群上做。第三條線才是模型——用訓(xùn)練腳本跑一遍邏輯回歸打開生成的模型報(bào)告確認(rèn)覆蓋率、區(qū)分度指標(biāo)在合理范圍。我最常踩的坑是把源碼里寫死的HDFS路徑直接拿來(lái)用。文檔里可能是/data/raw/credit_apply你的集群目錄結(jié)構(gòu)不一定一樣而且建表的LOCATION權(quán)限、Hive的warehouse路徑、Spark的checkpoint目錄每一處都要按實(shí)際環(huán)境改一遍。改完之后用一個(gè)小技巧驗(yàn)證完整性跑一次數(shù)據(jù)抽樣統(tǒng)計(jì)對(duì)比源文件的行數(shù)、字段數(shù)、空值率全部對(duì)上再繼續(xù)下一步。至少先斬釘截鐵跑通這個(gè)最小閉環(huán)你一晚上就能把整個(gè)系統(tǒng)的骨架摸清楚剩下的就是把規(guī)則閾值、模型參數(shù)按你的業(yè)務(wù)數(shù)據(jù)重新調(diào)一輪。希望幫到你。本文還有配套的精品資源點(diǎn)擊獲取