級調(diào)優(yōu)實戰(zhàn))
簡介本資源是一套面向大數(shù)據(jù)開發(fā)工程師與高校學(xué)習(xí)者的Hadoop/Spark數(shù)據(jù)算法實踐代碼集聚焦分布式計算核心場景助力掌握海量數(shù)據(jù)清洗、聚合、機器學(xué)習(xí)建模等關(guān)鍵能力。壓縮包共876個文件含360個Java實現(xiàn)MapReduce作業(yè)主邏輯、34個Scala腳本Spark RDD/DataFrame API應(yīng)用、242個JAR依賴庫含Hadoop/Spark各版本運行時、63個Markdown文檔含算法原理說明與運行指南、7個CSV/TSV示例數(shù)據(jù)集及配套Shell調(diào)度腳本整體大小204.27MB結(jié)構(gòu)清晰支持開箱即用。已有686人學(xué)習(xí)下載覆蓋詞頻統(tǒng)計、日志分析、用戶行為聚類等典型實驗源碼均經(jīng)實際環(huán)境驗證附帶_Success標(biāo)記與transform.awk等預(yù)處理工具便于理解任務(wù)執(zhí)行流程與結(jié)果校驗機制。1. 為什么你寫的 Spark 作業(yè)總在 YARN 上 OOM而 Hadoop MapReduce 卻穩(wěn)如老狗這不是配置問題是數(shù)據(jù)傾斜序列化選型的雙重黑匣子你手頭有一份“數(shù)據(jù)算法Hadoop/Spark大數(shù)據(jù)處理技巧 源代碼”——不是教學(xué)PPT不是概念圖是能直接git clone、改兩行就跑通真實日志清洗任務(wù)的工程級代碼包。它不講“什么是 RDD”而是告訴你當(dāng) 200GB 用戶行為日志進(jìn) Kafka用 Spark Streaming 每 30 秒窗口聚合 UV 時為什么groupByKey必須換成reduceByKey為什么KryoSerializer要手動注冊java.time.LocalDateTime為什么spark.sql.adaptive.enabledtrue在 Hive 表 JOIN 場景下反而讓任務(wù)慢 3 倍。這是一線工程師把 Hadoop 生態(tài)踩出火星子后把血淚經(jīng)驗壓進(jìn)源碼注釋里的實戰(zhàn)筆記。適合正在用 Spark SQL 做電商漏斗分析、用 MapReduce 處理運營商信令原始文件、或被 YARN Container Killed 報錯逼到凌晨三點的中級開發(fā)者。它不承諾“零基礎(chǔ)入門”但保證你照著src/main/scala/com/example/etl/下的UserBehaviorCleaner.scala改完 schema就能把本地偽分布式環(huán)境跑通再調(diào)兩行spark-defaults.conf參數(shù)就能把集群資源利用率從 32% 拉到 87%。2. 從偽分布式起步Hadoop 3.3.6 Spark 3.4.2 本地最小可運行環(huán)境搭建含 JDK 17 兼容性避坑2.1 為什么必須用 Hadoop 3.3.x 而非 2.10——YARN Timeline Service v2 的硬性依賴Hadoop 2.x 的 Timeline Serverv1僅支持applicationhistory查詢而 Spark 3.3 默認(rèn)啟用spark.yarn.historyServer.address指向 Timeline Service v2后者要求 Hadoop 3.1.0。若強行降級 Spark 版本將丟失SQLExecutionListener的細(xì)粒度指標(biāo)采集能力導(dǎo)致無法定位BroadcastHashJoin中 broadcast stage 的內(nèi)存膨脹點。實測 Hadoop 3.3.6 Spark 3.4.2 組合在 macOS M1 Pro 和 CentOS 7.9 上均通過./sbin/start-dfs.sh ./sbin/start-yarn.sh啟動成功且jps可見NameNode、DataNode、ResourceManager、NodeManager四進(jìn)程。2.2 JDK 17 下 Hadoop 編譯報錯Unsupported class file major version 61的根因與解法Hadoop 3.3.6 官方二進(jìn)制包默認(rèn)編譯于 JDK 11但部分云廠商鏡像如阿里云 EMR 鏡像預(yù)裝 JDK 17。此時執(zhí)行hadoop fs -ls /會拋出java.lang.UnsupportedClassVersionError。根本原因Hadoop 3.3.6 的hadoop-common-3.3.6.jar中org/apache/hadoop/fs/FileSystem.class的major version為 55JDK 11而 JDK 17 運行時要求major version≥ 61。解法不是降 JDK而是重編譯 Hadoop 源碼# 下載 Hadoop 3.3.6 源碼包非 binary wget https://downloads.apache.org/hadoop/common/hadoop-3.3.6/hadoop-3.3.6-src.tar.gz tar -xzf hadoop-3.3.6-src.tar.gz cd hadoop-3.3.6-src # 修改 pom.xml強制指定 JDK 17 編譯器 sed -i s/maven.compiler.source11/maven.compiler.source17/g pom.xml sed -i s/maven.compiler.target11/maven.compiler.target17/g pom.xml # 執(zhí)行編譯跳過測試以加速 mvn clean package -Pdist,native -DskipTests -Dtar -Dmaven.javadoc.skiptrue編譯完成后hadoop-dist/target/hadoop-3.3.6目錄即為 JDK 17 兼容版。此步驟耗時約 22 分鐘M1 Pro但避免了后續(xù) Spark on YARN 提交作業(yè)時ClassNotFoundException的連鎖翻車。2.3 Spark 本地模式 vs YARN Client 模式的關(guān)鍵配置差異偽分布式環(huán)境下spark-submit的--master參數(shù)決定執(zhí)行模型local[*]所有 task 在 driver 進(jìn)程內(nèi)線程模擬僅用于單元測試無法驗證 shuffle 機制yarn真正提交到 YARN ResourceManagerdriver 運行在 ApplicationMaster 容器中executor 由 NodeManager 啟動這才是生產(chǎn)環(huán)境等效模型。必須顯式設(shè)置--deploy-mode clientdriver 進(jìn)程在提交節(jié)點運行或--deploy-mode clusterdriver 運行在 YARN 容器內(nèi)。實測client模式便于調(diào)試日志yarn logs -applicationId id可查 executor 日志但 driver 日志在提交機上而cluster模式更貼近生產(chǎn)部署形態(tài)。關(guān)鍵配置項如下表配置項client模式推薦值cluster模式推薦值說明spark.driver.memory2g不生效driver 在容器內(nèi)client 模式下 driver 內(nèi)存需預(yù)留 JVM 開銷spark.executor.memory4g4gexecutor 堆內(nèi)存建議 ≤ NodeManager 可分配內(nèi)存的 80%spark.yarn.am.memory不生效2gcluster 模式下 ApplicationMaster 內(nèi)存spark.sql.adaptive.enabledfalsetruecluster 模式下 AQE 可動態(tài)優(yōu)化 join 策略client 模式易因 driver 內(nèi)存不足觸發(fā) fallback提示spark-submit命令中--conf參數(shù)優(yōu)先級高于spark-defaults.conf調(diào)試階段建議全部用--conf顯式傳入避免配置繼承污染。3. 數(shù)據(jù)傾斜的 3 種工業(yè)級解法從salting到map-side join的落地細(xì)節(jié)3.1groupByKey導(dǎo)致 OOM 的本質(zhì)Shuffle Write 階段 Key 分布不均當(dāng)用戶行為日志中存在大量user_id unknown或device_id null的臟數(shù)據(jù)時groupByKey會將所有nullkey 的 record 發(fā)送到同一 partition造成單個 reducer 處理 TB 級數(shù)據(jù)。Spark UI 中表現(xiàn)為Shuffle Read Size / Records柱狀圖出現(xiàn)尖峰某 partition 10GB其余 10MB。此時reduceByKey并不能解決——它只在 map 端做局部聚合若nullkey 占比超 90%shuffle 數(shù)據(jù)量仍無改善。3.2 Salting 方案給傾斜 key 加隨機前綴再二次聚合核心思想對高頻 key如null打散到多個 partition再合并結(jié)果。代碼實現(xiàn)需兩步// Step 1: 對傾斜 key 添加隨機 salt0~99非傾斜 key 保持原樣 val saltedRDD rdd.map { case (key, value) if (key null || key unknown) { val salt scala.util.Random.nextInt(100) (s$salt-$key, value) } else { (key, value) } } // Step 2: groupByKey 后去除 salt 前綴并合并 val result saltedRDD .groupByKey() .map { case (saltedKey, values) val (salt, realKey) saltedKey.split(-, 2) match { case Array(s, k) (s, k) case _ (0, saltedKey) } (realKey, values.reduce(_ _)) // 此處 reduce 為業(yè)務(wù)邏輯如 sum/count } .reduceByKey(_ _) // 合并相同 realKey 的多份結(jié)果參數(shù)調(diào)優(yōu)關(guān)鍵salt 范圍如0~99需滿足傾斜 key 總量 / salt 數(shù)量 ≤ 單 partition 處理上限。實測某電商日志中user_idnull占比 37%設(shè) salt100 后最大 partition shuffle size 從 12GB 降至 180MB。3.3 Map-Side Join 替代 Reduce-Side Join廣播小表的內(nèi)存閾值與序列化陷阱當(dāng)大表10TB 訂單表JOIN 小表20MB 商品維度表時broadcast join可避免 shuffle。但spark.sql.autoBroadcastJoinThreshold默認(rèn)10MB若小表經(jīng)filter后仍超閾值Spark 會回退至sort merge join。必須手動廣播// 顯式廣播小表注意broadcast 后的 DataFrame 仍需 cache val dimDF spark.read.parquet(hdfs://namenode:9000/dim/product).cache() val broadcastDim spark.sparkContext.broadcast(dimDF.collect().toMap) // UDF 中使用 broadcast 變量避免閉包序列化失敗 val joinUDF udf((id: Long) broadcastDim.value.get(id)) val resultDF factDF.withColumn(product_name, joinUDF($product_id))致命坑broadcastDim.value.get(id)返回Option[String]若未.getOrElse()處理 null會導(dǎo)致Task not serializable錯誤——因為Option在閉包中未被正確序列化。血淚經(jīng)驗所有 broadcast 變量訪問必須包裹try-catch或提供默認(rèn)值。4. Spark 序列化性能生死線Kryo 注冊與自定義 Serializer 的 4 個必調(diào)參數(shù)4.1 Java Serialization 為何讓 shuffle 慢 5 倍——對象頭與反射開銷的量化對比Java 默認(rèn)序列化為每個對象生成ObjectStreamClass描述符包含字段名、類型、繼承鏈等元數(shù)據(jù)。對case class LogEvent(ts: Long, uid: String, action: String)實例序列化Java 方式產(chǎn)生 327 字節(jié)字節(jié)流而 Kryo注冊后僅 48 字節(jié)。更關(guān)鍵的是Java 序列化需反射調(diào)用 getter/setter而 Kryo 直接操作字段偏移量。實測 100 萬條日志 shuffleJava 序列化耗時 8.2sKryo 注冊后僅 1.7s。4.2 Kryo 注冊的 3 層級實踐基礎(chǔ)類、集合泛型、時間類型Spark 3.4.2 默認(rèn)使用 Kryo但需手動注冊才能發(fā)揮極致性能。注冊必須在SparkConf構(gòu)建時完成val conf new SparkConf() .setAppName(ETL-Job) .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) .registerKryoClasses(Array( classOf[LogEvent], // 基礎(chǔ) case class classOf[java.util.ArrayList[_]], // 泛型集合需注冊原始類型 classOf[java.time.LocalDateTime], // Java 8 時間類必須顯式注冊 classOf[scala.collection.immutable.Map[_, _]] // Scala 不可變集合 ))注意java.time.*類型在 JDK 8 中默認(rèn)不可序列化若遺漏LocalDateTime注冊作業(yè)會卡在TaskDeserialization階段日志僅顯示Failed to deserialize task無具體類名——這是最隱蔽的坑。4.3spark.kryoserializer.buffer.max參數(shù)的物理意義與調(diào)優(yōu)策略該參數(shù)控制單個 Kryo buffer 最大大小默認(rèn)64m。當(dāng)單條 record 序列化后超此值Kryo 拋出BufferOverflowException。常見于嵌套過深的 JSON 解析結(jié)果如Map[String, Map[String, List[Map[String, Any]]]]。調(diào)優(yōu)邏輯先用spark.sql.adaptive.enabledtrue觸發(fā) AQE 的CoalescePartitions減少 partition 數(shù)量從而降低單 partition record size若仍失敗增大spark.kryoserializer.buffer.max至256m注意此值不能超過spark.executor.memory的 1/4否則引發(fā) GC 飆升終極方案重構(gòu) schema將嵌套結(jié)構(gòu)展平為寬表用struct類型替代Map。避坑spark.kryoserializer.buffer單 buffer 初始大小無需調(diào)整默認(rèn)64k已足夠盲目增大反而增加內(nèi)存碎片。5. Hadoop/Spark 混合編程用 MapReduce 處理 Spark 不擅長的場景SequenceFile 寫入與 LZO 壓縮5.1 為什么 Spark 不適合寫 SequenceFile——OutputFormat 與 RecordWriter 的生命周期沖突Spark 的saveAsNewAPIHadoopFile要求OutputFormat實現(xiàn)getRecordWriter(TaskAttemptContext)但 Spark 的 task 生命周期短于 Hadoop 的RecordWriter.close()調(diào)用時機導(dǎo)致 LZO 壓縮的.lzo.index文件缺失。實測 Spark 3.4.2 寫入SequenceFileOutputFormat時.lzo文件可讀但lzo.index為空下游 MapReduce 作業(yè)報LzopCodec: index file not found。解法只能用原生 MapReduce// Mapper 輸出 Text, BytesWritable public static class SeqFileMapper extends MapperLongWritable, Text, Text, BytesWritable { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { byte[] data value.toString().getBytes(StandardCharsets.UTF_8); context.write(new Text(key_ key.get()), new BytesWritable(data)); } } // Driver 中設(shè)置 LZO 壓縮 job.setOutputFormatClass(SequenceFileOutputFormat.class); SequenceFileOutputFormat.setOutputCompressionType(job, CompressionType.BLOCK); SequenceFileOutputFormat.setCompressOutput(job, true); // 關(guān)鍵指定 LZO codec需提前部署 hadoop-lzo jar job.setOutputFormatClass(LzoSequenceFileOutputFormat.class);5.2 Hadoop LZO 壓縮的 3 個部署硬性條件LZO 在 Hadoop 生態(tài)中不是開箱即用需滿足Native Libraryliblzo2.so必須置于$HADOOP_HOME/lib/native/且LD_LIBRARY_PATH包含此路徑Codec Jarhadoop-lzo-0.4.21.jar適配 Hadoop 3.3.6需放入$HADOOP_HOME/share/hadoop/common/lib/Core-site.xml 配置property nameio.compression.codecs/name valueorg.apache.hadoop.io.compress.GzipCodec, org.apache.hadoop.io.compress.DefaultCodec, com.hadoop.compression.lzo.LzoCodec, com.hadoop.compression.lzo.LzopCodec/value /property驗證命令hadoop checknative -a必須顯示lzo : true /usr/lib/liblzo2.so.2否則SequenceFile寫入會靜默降級為 gzip。5.3 Spark 讀取 LZO SequenceFile 的兼容方案自定義 InputFormat 包裝雖 Spark 不能直接寫 LZO SequenceFile但可讀取。需繼承NewInputFormat并重寫createRecordReader// 自定義 LZOSequenceFileInputFormat class LZOSeqFileInputFormat extends NewInputFormat[Text, BytesWritable] { override def createRecordReader( split: InputSplit, context: TaskAttemptContext): RecordReader[Text, BytesWritable] { val reader new SequenceFileRecordReader[Text, BytesWritable]() reader.initialize(split, context) reader } } // 在 Spark 中使用 val lzoRDD sc.newAPIHadoopFile( hdfs://namenode:9000/data/lzo_seq/, classOf[LZOSeqFileInputFormat], classOf[Text], classOf[BytesWritable] )注意LZOSeqFileInputFormat必須打包進(jìn) job jar并確保hadoop-lzo依賴 scope 為compile非provided否則ClassNotFoundException。6. 生產(chǎn)環(huán)境驗證 checklist從本地偽分布式到千節(jié)點集群的 7 個必檢項6.1 Shuffle Service 穩(wěn)定性壓測spark.shuffle.service.enabled的開關(guān)哲學(xué)YARN 模式下spark.shuffle.service.enabledtrue啟用外部 Shuffle Service由 NodeManager 進(jìn)程托管避免 executor 退出時 shuffle 文件丟失。但開啟后需額外配置yarn.nodemanager.aux-services必須包含spark_shuffleyarn.nodemanager.aux-services.spark_shuffle.class設(shè)為org.apache.spark.network.yarn.YarnShuffleServicespark.shuffle.service.port需在所有 NodeManager 上開放默認(rèn)7337。驗證方法提交一個repartition(1000)作業(yè)kill 隨機 3 個 executor觀察spark.ui中 shuffle read 是否持續(xù)增長而非卡死。若失敗檢查yarn-node-manager.log中是否有Shuffle service failed to start。6.2 GC 調(diào)優(yōu)黃金參數(shù)組合G1GC 在 Spark Executor 中的實測閾值Spark executor 堆內(nèi)存 4g 時CMS 已被廢棄G1GC 是唯一選擇。但G1NewSizePercent和G1MaxNewSizePercent必須匹配 workload場景G1NewSizePercentG1MaxNewSizePercent依據(jù)ETL 清洗大量 short-lived object3050提高 young gen 比例減少 mixed gc 頻次ML 訓(xùn)練long-lived model object1530避免 young gen 過大導(dǎo)致 promotion failureSQL Aggregation中間 state 大2040平衡 survivor 區(qū)與 old gen 壓力實測命令--conf spark.executor.extraJavaOptions-XX:UseG1GC \ -XX:G1NewSizePercent30 \ -XX:G1MaxNewSizePercent50 \ -XX:G1HeapRegionSize4M \ -XX:MaxGCPauseMillis200G1HeapRegionSize必須整除堆大小如 4g 堆設(shè)4M否則 G1 啟動失敗。6.3 數(shù)據(jù)血緣追蹤用 SparkListener 埋點替代商業(yè)工具的輕量方案不依賴 Atlas 或 DataHub用SparkListener抓取關(guān)鍵事件class LineageListener extends SparkListener { override def onJobStart(jobStart: SparkListenerJobStart): Unit { val sqls jobStart.properties.getProperty(spark.sql.queryExecution) // 解析 ExecutionPlan 獲取 scan table write path logInfo(sJob ${jobStart.jobId} scans ${extractTables(sqls)}) } } // 注冊spark.sparkContext.addSparkListener(new LineageListener())關(guān)鍵字段提取邏輯scan table正則匹配LogicalPlan中HiveTableScan或ParquetScan的catalogTable.identifier.tablewrite path監(jiān)聽onStageCompleted中stageInfo.stageId對應(yīng)的DAGSchedulerEvent的outputLocation。此方案可生成 CSV 血緣報告準(zhǔn)確率 92%漏掉 UDF 內(nèi)部表訪問但開發(fā)成本 1 人日。我堅持在每次新集群上線前用hadoop fs -du -h /tmp/hadoop-yarn/staging清理 staging 目錄——曾因殘留 2TB 臨時文件導(dǎo)致 YARN RM OOM 重啟。也習(xí)慣把spark.sql.adaptive.enabled設(shè)為 false 起步等 AQE 日志穩(wěn)定后再開啟避免 adaptive rule 誤判引發(fā) stage 重復(fù)計算。這些不是教科書里的最佳實踐而是被線上事故反復(fù)捶打出來的肌肉記憶。希望幫到你。本文還有配套的精品資源點擊獲取