戰(zhàn):從集群搭建到 OOM 排查與調(diào)優(yōu)指南)
最近圈子里討論度比較高的一個(gè)話題是Muse Spark周榜沖入前三。很多人第一反應(yīng)是這到底是個(gè)新框架還是某個(gè)團(tuán)隊(duì)內(nèi)部項(xiàng)目的名字其實(shí)嚴(yán)格來(lái)說(shuō)Muse Spark 本身并不是一個(gè)獨(dú)立的大數(shù)據(jù)計(jì)算引擎真正支撐它上榜的是背后那套已經(jīng)在大數(shù)據(jù)領(lǐng)域沉淀了十多年的 Apache Spark 技術(shù)棧。這個(gè)話題能沖上熱榜恰恰說(shuō)明了一件事在 2025 年這個(gè)時(shí)間點(diǎn)Apache Spark 依然是大數(shù)據(jù)分析、數(shù)據(jù)倉(cāng)庫(kù)、實(shí)時(shí)計(jì)算場(chǎng)景里繞不開(kāi)的底座。無(wú)論是做離線 ETL、即席查詢、機(jī)器學(xué)習(xí)特征工程還是跑流式數(shù)據(jù)處理Spark 都是面試高頻、生產(chǎn)高頻、踩坑也高頻的技術(shù)方向。這篇文章不會(huì)去重復(fù)官方文檔里那些概念定義而是從工程落地的角度把 Spark 從環(huán)境搭建、集群配置、SQL 數(shù)據(jù)分析、外部數(shù)據(jù)源集成到 OOM 排查、性能調(diào)優(yōu)、信創(chuàng)數(shù)據(jù)庫(kù)適配這些真實(shí)場(chǎng)景完整過(guò)一遍。如果你正準(zhǔn)備搭建一套 Spark 環(huán)境、正在處理 Spark SQL 的性能問(wèn)題、或者馬上要面大數(shù)據(jù)崗位這篇文章值得收藏備用。1. 從周榜沖前三看 Spark 的長(zhǎng)期價(jià)值先聊聊為什么 Spark 相關(guān)的關(guān)鍵詞能持續(xù)保持熱度。一個(gè)很重要的原因是Spark 的生態(tài)位太特殊了它的底層是分布式計(jì)算引擎上層卻能同時(shí)覆蓋 SQL、流處理、機(jī)器學(xué)習(xí)和圖計(jì)算。這意味著一個(gè)團(tuán)隊(duì)只要引入 Spark就能統(tǒng)一處理離線批任務(wù)、實(shí)時(shí)流任務(wù)和部分 AI 訓(xùn)練前的數(shù)據(jù)處理工作技術(shù)棧能省掉好幾套。從實(shí)際招聘需求來(lái)看Spark 相關(guān)崗位的面試題也非常固定Spark 任務(wù)提交流程、寬依賴窄依賴、Shuffle 原理、數(shù)據(jù)傾斜、內(nèi)存調(diào)優(yōu)、Spark SQL 執(zhí)行計(jì)劃……這些問(wèn)題不僅面試會(huì)考生產(chǎn)環(huán)境里幾乎每天都會(huì)遇到。熱搜詞里的spark oom、spark sql、spark 讀取 redis、spark 集群搭建基本就是一線工程師最常搜索的排障和開(kāi)發(fā)關(guān)鍵詞。所以與其糾結(jié) Muse Spark 這個(gè)名詞本身不如把它看作一個(gè)信號(hào)Spark 技術(shù)棧依然是當(dāng)前數(shù)據(jù)工程領(lǐng)域最值得投資的學(xué)習(xí)方向。這篇文章的核心目標(biāo)就是把這些高頻問(wèn)題串起來(lái)給出一條從零到一、再到生產(chǎn)可用的完整路徑。2. Spark 核心概念與適用場(chǎng)景在動(dòng)手之前先把幾個(gè)核心概念理清楚。很多新手對(duì) Spark 的誤解往往集中在 RDD、DataFrame、DataSet 這三者的關(guān)系上以及 Spark 和 Hadoop MapReduce 的差異上。2.1 Spark 到底是什么Spark 是一個(gè)基于內(nèi)存計(jì)算的分布式數(shù)據(jù)處理引擎。它的核心思想是把一份大數(shù)據(jù)拆成多個(gè)分區(qū)分發(fā)到集群的多臺(tái)機(jī)器上并行計(jì)算然后在必要時(shí)對(duì)數(shù)據(jù)進(jìn)行重新分區(qū)Shuffle最終匯總結(jié)果。對(duì)比 Hadoop MapReduceSpark 最大的優(yōu)勢(shì)在于中間結(jié)果可以緩存在內(nèi)存里而 MapReduce 每一步幾乎都要落盤(pán)所以 Spark 在迭代計(jì)算、交互式查詢、機(jī)器學(xué)習(xí)這類場(chǎng)景下性能往往高出數(shù)倍甚至數(shù)十倍。代價(jià)是 Spark 對(duì)內(nèi)存資源的規(guī)劃要求更高這也是 OOM 問(wèn)題高發(fā)的根本原因。2.2 RDD、DataFrame、DataSet 的區(qū)別概念特點(diǎn)適用場(chǎng)景RDD底層抽象彈性分布式數(shù)據(jù)集適合非結(jié)構(gòu)化數(shù)據(jù)和自定義計(jì)算邏輯需要精細(xì)控制分區(qū)、依賴關(guān)系時(shí)DataFrame帶 Schema 的分布式表有列名和類型底層基于 Row大多數(shù)數(shù)據(jù)分析、SQL 場(chǎng)景最常用DataSet強(qiáng)類型支持 Java/Scala 類型安全Python 中無(wú)對(duì)應(yīng)JVM 系團(tuán)隊(duì)追求編譯期類型檢查時(shí)實(shí)際開(kāi)發(fā)中90% 以上的場(chǎng)景用 DataFrame 就夠了。RDD 更多出現(xiàn)在底層框架開(kāi)發(fā)和非常復(fù)雜的自定義算子中。這里真正容易踩坑的地方是一些同學(xué)在 PySpark 里強(qiáng)行使用 RDD 寫(xiě) map 算子結(jié)果丟失了 Catalyst 優(yōu)化器的優(yōu)化能力同樣的數(shù)據(jù)量性能差距可能十倍以上。2.3 適用場(chǎng)景判斷Spark 適合的場(chǎng)景包括離線 ETL 清洗、大規(guī)模日志分析、復(fù)雜 SQL 查詢、特征工程、流式數(shù)據(jù)處理Structured Streaming、圖計(jì)算。不太適合的場(chǎng)景是低延遲毫秒級(jí)在線查詢、單表幾十萬(wàn)行以內(nèi)的小數(shù)據(jù)分析。如果你只有幾十萬(wàn)行數(shù)據(jù)用 ClickHouse、Doris甚至關(guān)系型數(shù)據(jù)庫(kù)反而更快強(qiáng)行上 Spark 只會(huì)增加運(yùn)維成本。3. Spark 環(huán)境準(zhǔn)備與集群搭建無(wú)論是學(xué)習(xí)還是生產(chǎn)第一步都是把環(huán)境跑起來(lái)。這里以 Linux 環(huán)境下的 Standalone 集群為例演示最經(jīng)典的搭建方式。整體思路同樣適用于 YARN、Kubernetes 模式只是資源調(diào)度層不同。3.1 前置條件操作系統(tǒng)LinuxCentOS 7 / Ubuntu 20.04 均可JDK1.8 或 11Spark 3.x 版本建議 JDK 8/11具體以官方說(shuō)明為準(zhǔn)Python如果使用 PySpark3.8 以上SSH 免密鑰登錄集群節(jié)點(diǎn)之間需要互通需要說(shuō)明的是本文不寫(xiě)死具體版本號(hào)因?yàn)?Spark 版本更新較快建議到官網(wǎng)下載當(dāng)前穩(wěn)定版重點(diǎn)理解配置思路。3.2 安裝包準(zhǔn)備# 下載 Spark 安裝包選擇 pre-built for Apache Hadoop 版本 tar -zxvf spark-*.tgz mv spark-* /opt/spark # 配置環(huán)境變量 cat ~/.bashrc EOF export SPARK_HOME/opt/spark export PATH$SPARK_HOME/bin:$SPARK_HOME/sbin:$PATH export PYSPARK_PYTHONpython3 EOF source ~/.bashrc # 驗(yàn)證安裝 spark-shell --version3.3 配置 Standalone 集群cd /opt/spark/conf # 1. 修改 spark-env.sh cp spark-env.sh.template spark-env.sh cat spark-env.sh EOF JAVA_HOME/opt/jdk SPARK_MASTER_HOSTnode01 SPARK_MASTER_PORT7077 SPARK_WORKER_CORES4 SPARK_WORKER_MEMORY8g SPARK_WORKER_INSTANCES1 EOF # 2. 配置 workers cp workers.template workers cat workers EOF node01 node02 node03 EOF # 3. 配置 spark-defaults.conf cp spark-defaults.conf.template spark-defaults.conf cat spark-defaults.conf EOF spark.master spark://node01:7077 spark.serializer org.apache.spark.serializer.KryoSerializer spark.sql.shuffle.partitions 8 spark.driver.memory 2g spark.executor.memory 4g spark.executor.cores 2 EOF3.4 啟動(dòng)與驗(yàn)證# 啟動(dòng)集群 /opt/spark/sbin/start-master.sh /opt/spark/sbin/start-workers.sh # 查看進(jìn)程 jps # Master 進(jìn)程會(huì)顯示 Master # 每個(gè) Worker 節(jié)點(diǎn)會(huì)顯示 Worker # 訪問(wèn) Web UI # http://node01:8080啟動(dòng)后可以在 Web UI 上看到 Worker 節(jié)點(diǎn)的資源情況總內(nèi)存、總核數(shù)、已用資源。這里真正容易踩坑的地方是SPARK_WORKER_MEMORY和spark.executor.memory的概念不能混。前者是 Worker 進(jìn)程給這個(gè)節(jié)點(diǎn)上所有 Executor 能用的總內(nèi)存上限后者是每個(gè) Executor 的內(nèi)存大小。一個(gè) Worker 節(jié)點(diǎn)上如果跑了多個(gè) Executor累加值不能超過(guò) Worker 總內(nèi)存否則任務(wù)會(huì)一直卡在等待資源的狀態(tài)。4. Spark SQL 數(shù)據(jù)分析完整示例集群搭好之后用一個(gè)最典型的數(shù)據(jù)分析場(chǎng)景來(lái)跑通流程讀取一份商品訂單 CSV 數(shù)據(jù)做分組聚合統(tǒng)計(jì)。這一步能同時(shí)驗(yàn)證集群、Spark SQL、執(zhí)行計(jì)劃這三層是否正常。4.1 準(zhǔn)備測(cè)試數(shù)據(jù)mkdir -p /data/spark-demo cat /data/spark-demo/orders.csv EOF order_id,user_id,category,amount,order_date 1001,U001,手機(jī),2999,2025-01-01 1002,U002,電腦,5999,2025-01-01 1003,U001,手機(jī),3499,2025-01-02 1004,U003,家電,1299,2025-01-02 1005,U002,電腦,7999,2025-01-03 EOF4.2 編寫(xiě) PySpark 分析腳本# 文件路徑/data/spark-demo/order_analysis.py from pyspark.sql import SparkSession from pyspark.sql.functions import sum, count, avg # 創(chuàng)建 SparkSession spark SparkSession.builder \ .appName(OrderAnalysis) \ .getOrCreate() # 讀取 CSV 數(shù)據(jù)自動(dòng)推斷 Schema df spark.read \ .option(header, true) \ .option(inferSchema, true) \ .csv(/data/spark-demo/orders.csv) # 查看 Schema print( Schema ) df.printSchema() # 注冊(cè)臨時(shí)視圖便于寫(xiě) SQL df.createOrReplaceTempView(orders) # 場(chǎng)景1按品類統(tǒng)計(jì)訂單量、總金額、平均金額 print( 品類聚合統(tǒng)計(jì) ) result1 spark.sql( SELECT category, COUNT(order_id) AS order_cnt, SUM(amount) AS total_amount, AVG(amount) AS avg_amount FROM orders GROUP BY category ORDER BY total_amount DESC ) result1.show() # 場(chǎng)景2統(tǒng)計(jì)每個(gè)用戶的消費(fèi)總金額與訂單數(shù) print( 用戶消費(fèi)統(tǒng)計(jì) ) result2 df.groupBy(user_id) \ .agg( count(order_id).alias(order_cnt), sum(amount).alias(total_spent) ) \ .orderBy(total_spent, ascendingFalse) result2.show() # 停止 SparkSession spark.stop()這段腳本包含了兩個(gè)最核心的 Spark SQL 開(kāi)發(fā)方式一是直接寫(xiě) SQL 字符串適合從 Hive SQL 遷過(guò)來(lái)的團(tuán)隊(duì)二是使用 DataFrame API適合在代碼里做復(fù)雜的條件邏輯。兩種方式最終都會(huì)經(jīng)過(guò) Catalyst 優(yōu)化器執(zhí)行效率沒(méi)有本質(zhì)區(qū)別選擇哪種主要看團(tuán)隊(duì)的編碼習(xí)慣。4.3 提交任務(wù)/opt/spark/bin/spark-submit \ --master spark://node01:7077 \ --deploy-mode client \ /data/spark-demo/order_analysis.py4.4 運(yùn)行結(jié)果驗(yàn)證正常情況下控制臺(tái)會(huì)依次輸出 Schema、品類聚合結(jié)果、用戶消費(fèi)結(jié)果。以品類統(tǒng)計(jì)為例預(yù)期輸出大致如下--------------------------------------- |category|order_cnt|total_amount|avg_amount| --------------------------------------- |電腦 |2 |13998.0 |6999.0 | |手機(jī) |2 |6498.0 |3249.0 | |家電 |1 |1299.0 |1299.0 | ---------------------------------------如何判斷成功三個(gè)標(biāo)準(zhǔn)任務(wù)退出碼為 0沒(méi)有拋出 Exception。輸出結(jié)果與手工計(jì)算一致。Spark Web UI端口 4040上能看到完整的 Job、Stage、Task 記錄并且沒(méi)有大量失敗重試。如果失敗第一步先看 Executor 日志中的ERROR信息而不是只看 Driver 端的告警。大多數(shù) Spark SQL 任務(wù)失敗的真實(shí)原因都藏在 Executor 日志里。5. Spark 讀取外部數(shù)據(jù)源以 Redis 為例生產(chǎn)環(huán)境中Spark 很少只讀取 HDFS 或本地文件最常見(jiàn)的其實(shí)是和 Redis、Kafka、MySQL、達(dá)夢(mèng)等外部系統(tǒng)對(duì)接。這里以spark 讀取 redis為例講清楚外部數(shù)據(jù)源的接入套路。5.1 為什么需要 Spark 讀 Redis一個(gè)典型的場(chǎng)景是有一批離線任務(wù)需要基于 Redis 中的黑白名單、實(shí)時(shí)指標(biāo)緩存或維度數(shù)據(jù)進(jìn)行關(guān)聯(lián)分析。如果在 Spark 任務(wù)里逐條調(diào)用 Redis 客戶端會(huì)產(chǎn)生大量網(wǎng)絡(luò)請(qǐng)求性能極差。正確做法是使用官方提供的 Redis-Spark 連接器利用 RDD 分區(qū)并行讀取。5.2 添加依賴如果使用 Maven 管理 Scala/Java 項(xiàng)目需要引入dependency groupIdcom.redislabs/groupId artifactIdspark-redis/artifactId version版本號(hào)以官方 Maven 倉(cāng)庫(kù)為準(zhǔn)/version /dependency使用 PySpark 時(shí)需要把對(duì)應(yīng)的 JAR 包通過(guò)--jars參數(shù)傳給 spark-submit。5.3 代碼示例# 文件路徑/data/spark-demo/redis_demo.py from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(RedisDemo) \ .config(spark.redis.host, 192.168.1.100) \ .config(spark.redis.port, 6379) \ .config(spark.redis.auth, 你的密碼) \ .getOrCreate() # 讀取 Redis 中的 Hash 類型數(shù)據(jù) df spark.read \ .format(org.apache.spark.sql.redis) \ .option(table, user_profile) \ .load() df.show(10) spark.stop()需要特別提醒的是讀取 Redis 時(shí)最好加上filter條件或者明確限制分區(qū)數(shù)。Redis 本身不是為全量掃描設(shè)計(jì)的如果整個(gè) Key 空間都導(dǎo)進(jìn) SparkRedis 服務(wù)端很可能會(huì)成為瓶頸甚至拖垮在線業(yè)務(wù)。生產(chǎn)環(huán)境更常見(jiàn)的方案是只讀取特定前綴的 Key或者在離線低峰期執(zhí)行。6. Spark SQL 性能優(yōu)化與 OOM 排查熱搜詞里spark oom排名非常靠前這確實(shí)是生產(chǎn)環(huán)境中最常見(jiàn)的問(wèn)題之一。Spark OOM 不是單一原因需要分場(chǎng)景排查。6.1 常見(jiàn)的 OOM 類型現(xiàn)象可能原因典型特征Driver OOMcollect() 拉取數(shù)據(jù)量過(guò)大日志報(bào) java.lang.OutOfMemoryError: Java heap spaceExecutor OOM(執(zhí)行內(nèi)存)單個(gè)任務(wù)數(shù)據(jù)量過(guò)大或數(shù)據(jù)傾斜Shuffle Read 階段報(bào)內(nèi)存溢出Executor OOM(存儲(chǔ)內(nèi)存)cache/persist 數(shù)據(jù)超過(guò)存儲(chǔ)上限Storage 頁(yè)簽顯示內(nèi)存溢出堆外內(nèi)存溢出網(wǎng)絡(luò)讀取、序列化數(shù)據(jù)量過(guò)大Direct buffer memory 異常6.2 最核心的優(yōu)化手段先說(shuō)一個(gè)容易被忽略的結(jié)論Spark OOM 的第一排查順序不是調(diào)大內(nèi)存而是檢查數(shù)據(jù)分布是否傾斜。數(shù)據(jù)傾斜會(huì)導(dǎo)致某個(gè) Task 拉取了遠(yuǎn)超平均量的數(shù)據(jù)直接把 Executor 內(nèi)存打爆。這時(shí)候調(diào)大 Executor 內(nèi)存只是治標(biāo)甚至可能把節(jié)點(diǎn)物理內(nèi)存吃滿。# 提交任務(wù)時(shí)加入以下參數(shù)進(jìn)行調(diào)優(yōu) spark-submit \ --master spark://node01:7077 \ --executor-memory 8g \ --conf spark.sql.shuffle.partitions200 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ /data/spark-demo/order_analysis.pySpark 3.x 開(kāi)啟了 AQE自適應(yīng)查詢執(zhí)行之后很多原本需要手工調(diào)優(yōu)的參數(shù)可以自動(dòng)處理。skewJoin.enabledtrue能夠自動(dòng)拆分傾斜的分區(qū)這是目前解決數(shù)據(jù)傾斜性價(jià)比最高的手段。6.3 數(shù)據(jù)傾斜的定位方法在 Spark Web UI 的 Stage 頁(yè)簽里看某個(gè) Stage 的 Task 耗時(shí)和 Shuffle Read 數(shù)據(jù)量。如果發(fā)現(xiàn)少量 Task 處理的數(shù)據(jù)量是其他 Task 的幾倍甚至幾十倍基本可以斷定存在數(shù)據(jù)傾斜。解決方案除了 AQE還有加隨機(jī)前綴打散、廣播小表、兩階段聚合先局部聚合再加鹽后全局聚合等方法。7. 達(dá)夢(mèng)數(shù)據(jù)庫(kù)與 Spark 的適配集成達(dá)夢(mèng)數(shù)據(jù)庫(kù)(dm) 與 apache spark 的適配集成能成為熱搜詞說(shuō)明不少團(tuán)隊(duì)正在做信創(chuàng)環(huán)境下的數(shù)據(jù)遷移和平臺(tái)建設(shè)。達(dá)夢(mèng)是國(guó)產(chǎn)關(guān)系型數(shù)據(jù)庫(kù)在很多政企項(xiàng)目中承擔(dān)著核心數(shù)據(jù)庫(kù)的角色。Spark 和達(dá)夢(mèng)的集成本質(zhì)上是 Spark 通過(guò) JDBC 讀寫(xiě)達(dá)夢(mèng)數(shù)據(jù)庫(kù)。7.1 集成思路Spark 本身并不關(guān)心底層是什么數(shù)據(jù)庫(kù)只要提供了 JDBC 驅(qū)動(dòng)就能通過(guò)format(jdbc)的方式讀寫(xiě)。因此適配達(dá)夢(mèng)的核心步驟只有兩步拿到達(dá)夢(mèng)的 JDBC 驅(qū)動(dòng)配置正確的 JDBC URL。7.2 通過(guò) JDBC 讀取達(dá)夢(mèng)數(shù)據(jù)# 文件路徑/data/spark-demo/dm_demo.py from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(DMIntegration) \ .getOrCreate() # 讀取達(dá)夢(mèng)數(shù)據(jù)庫(kù)中的表 jdbc_df spark.read \ .format(jdbc) \ .option(url, jdbc:dm://192.168.1.101:5236) \ .option(user, test_user) \ .option(password, 你的密碼) \ .option(dbtable, schema_name.test_table) \ .option(driver, dm.jdbc.driver.DmDriver) \ .option(fetchsize, 1000) \ .load() # 統(tǒng)計(jì)行數(shù)驗(yàn)證連通性 print(行數(shù)統(tǒng)計(jì):, jdbc_df.count()) jdbc_df.show(5) spark.stop()7.3 寫(xiě)入達(dá)夢(mèng)數(shù)據(jù)庫(kù)# 提交任務(wù)時(shí)需要把達(dá)夢(mèng) JDBC 驅(qū)動(dòng) JAR 放到 --jars 中 spark-submit \ --master spark://node01:7077 \ --jars /opt/dm/dmsjdbc.jar \ /data/spark-demo/dm_demo.py寫(xiě)入時(shí)建議使用mode(overwrite)或mode(append)并且注意控制寫(xiě)入并發(fā)。達(dá)夢(mèng)作為關(guān)系型數(shù)據(jù)庫(kù)寫(xiě)入吞吐能力遠(yuǎn)不如大數(shù)據(jù)存儲(chǔ)系統(tǒng)如果 Executor 并發(fā)過(guò)高反而會(huì)把數(shù)據(jù)庫(kù)打死。生產(chǎn)環(huán)境建議通過(guò)repartition(10)限制寫(xiě)入并行度。需要說(shuō)明的是不同版本的達(dá)夢(mèng)數(shù)據(jù)庫(kù)驅(qū)動(dòng) JAR 包名可能不同具體以達(dá)夢(mèng)官方提供的驅(qū)動(dòng)為準(zhǔn)。如果遇到連接失敗重點(diǎn)檢查端口、用戶名權(quán)限、驅(qū)動(dòng)類名這三個(gè)點(diǎn)而不要懷疑 Spark 本身。8. Spark 常見(jiàn)問(wèn)題與排查思路這里整理一份高頻問(wèn)題排查表覆蓋日常開(kāi)發(fā)中 80% 以上的報(bào)錯(cuò)場(chǎng)景。問(wèn)題現(xiàn)象可能原因排查方式解決方案任務(wù)一直處于 Accepted/WaitingExecutor 資源不足Web UI 查看可用內(nèi)存和核數(shù)調(diào)大 Executor 內(nèi)存或申請(qǐng)更多 Worker 資源運(yùn)行時(shí)報(bào) java.lang.OutOfMemoryErrorExecutor 內(nèi)存不足或數(shù)據(jù)傾斜查看 Executor 日志、Stage 數(shù)據(jù)分布調(diào)整 AQE 參數(shù)、增大 Executor 內(nèi)存、優(yōu)化數(shù)據(jù)分布報(bào) java.lang.ClassNotFoundException缺少第三方 JAR 包查看棧頂異常類名通過(guò) --jars 提交依賴 JAR報(bào) org.apache.spark.shuffle.FetchFailedException網(wǎng)絡(luò)抖動(dòng)或 Executor 崩潰查看失敗 Task 所在節(jié)點(diǎn)檢查網(wǎng)絡(luò)、重啟任務(wù)降低單 Task 數(shù)據(jù)量Spark SQL 查詢速度極慢沒(méi)有謂詞下推、Shuffle 數(shù)據(jù)量大查看 Spark UI 中的 Exchange 節(jié)點(diǎn)優(yōu)化 SQL、過(guò)濾提前、開(kāi)啟 AQE報(bào) java.sql.SQLException: ORA-XXXX數(shù)據(jù)不符合目標(biāo)表約束查看具體 SQL 錯(cuò)誤碼根據(jù)錯(cuò)誤碼清理臟數(shù)據(jù)或修改目標(biāo)表結(jié)構(gòu)這里要強(qiáng)調(diào)一個(gè)工程習(xí)慣遇到 Spark 報(bào)錯(cuò)優(yōu)先去 Spark Web UI 的 Executors 和 Stage 頁(yè)面截圖留證再根據(jù)排查表確認(rèn)原因。很多 Spark 任務(wù)在提交端顯示的堆棧并不完整只有到了 Executor 日志層面才能看到真正的根因。9. Spark 面試高頻問(wèn)題盤(pán)點(diǎn)熱搜詞里有spark面試題正好把面試中最容易考的幾類問(wèn)題做一個(gè)簡(jiǎn)要梳理。這些問(wèn)題不是死記硬背就能答好的每一個(gè)背后都對(duì)應(yīng)對(duì)源碼或生產(chǎn)實(shí)踐的理解。9.1 基礎(chǔ)原理類Spark 任務(wù)的提交流程客戶端提交 - Driver 啟動(dòng) - 生成 DAG - 劃分 Stage - Task 分發(fā)到 Executor。寬依賴和窄依賴的區(qū)別窄依賴指父 RDD 的每個(gè)分區(qū)最多被子 RDD 的一個(gè)分區(qū)使用不需要 Shuffle寬依賴需要 Shuffle。Shuffle 的過(guò)程Map 端寫(xiě)入內(nèi)存緩沖區(qū) - 溢寫(xiě)磁盤(pán) - Reduce 端拉取合并。9.2 代碼實(shí)戰(zhàn)類repartition和coalesce的區(qū)別repartition 會(huì)觸發(fā) Shufflecoalesce 默認(rèn)不觸發(fā)但可能導(dǎo)致數(shù)據(jù)分布不均。cache和persist的區(qū)別cache 是 persist 的 MEMORY_ONLY 級(jí)別persist 可以指定 StorageLevel。map和mapPartitions的區(qū)別map 每條記錄處理一次mapPartitions 每個(gè)分區(qū)處理一次適合批量初始化連接等場(chǎng)景。9.3 調(diào)優(yōu)類數(shù)據(jù)傾斜怎么解決。Executor 內(nèi)存怎么劃分Spark 3.x 中Executor 內(nèi)存分為 Reserved、User Memory、Spark MemoryExecution 和 Storage 共享。小文件問(wèn)題怎么處理寫(xiě)入前 repartition/coalesce或者動(dòng)態(tài)分區(qū)裁剪。面試官真正想聽(tīng)的不是概念背誦而是你有沒(méi)有在生產(chǎn)環(huán)境踩過(guò)坑。比如數(shù)據(jù)傾斜問(wèn)題如果只是說(shuō)加隨機(jī)前綴而沒(méi)有說(shuō)具體怎么判斷傾斜、怎么驗(yàn)證效果面試官通常還會(huì)繼續(xù)追問(wèn)細(xì)節(jié)。10. 最佳實(shí)踐與工程建議最后結(jié)合這段時(shí)間的實(shí)戰(zhàn)經(jīng)驗(yàn)和近期熱榜上大家關(guān)心的問(wèn)題給出幾條建議。10.1 開(kāi)發(fā)階段一律使用 SparkSession不要再用舊的 SQLContext/HiveContext。SQL 和 DataFrame API 根據(jù)場(chǎng)景混用復(fù)雜邏輯用 SQL 表達(dá)更清晰動(dòng)態(tài)條件用 DataFrame API 更靈活。聚合結(jié)果先show()確認(rèn)不要直接collect()到 Driver避免 OOM。開(kāi)發(fā)環(huán)境加--conf spark.ui.port4040可以通過(guò) Web UI 實(shí)時(shí)觀察 Job 過(guò)程。10.2 生產(chǎn)階段核心任務(wù)一定要開(kāi)啟 AQESpark 3.2 之后 AQE 默認(rèn)開(kāi)啟但需要確認(rèn)集群版本。提交任務(wù)時(shí)通過(guò)spark-submit --conf傳參不要在代碼里硬編碼環(huán)境相關(guān)的 IP 和端口。外部數(shù)據(jù)源Redis、MySQL、達(dá)夢(mèng)的訪問(wèn)必須控制并發(fā)和超時(shí)時(shí)間防止影響在線服務(wù)。生產(chǎn)表和臨時(shí)表分開(kāi)命名避免同名覆蓋。使用spark.sql.shuffle.partitions時(shí)不要盲目調(diào)大。分區(qū)數(shù)越大的確越能分散壓力但每個(gè)分區(qū)都會(huì)產(chǎn)生對(duì)應(yīng)的 tasks調(diào)度開(kāi)銷同樣不可忽略通常設(shè)置為 Executor 總核數(shù)的 2~3 倍比較合理。10.3 團(tuán)隊(duì)協(xié)作Spark 任務(wù)腳本納入 Git 管理提交信息中注明業(yè)務(wù)口徑。建立任務(wù)血緣文檔說(shuō)明輸入表、輸出表、調(diào)度依賴和告警負(fù)責(zé)人。每次上線前先在測(cè)試環(huán)境跑通并用--executor-memory參數(shù)壓測(cè)一小部分?jǐn)?shù)據(jù)觀察資源水位。11. 總結(jié)從沖榜到真正的工程落地回到標(biāo)題。Muse Spark周榜沖進(jìn)前三不代表出現(xiàn)了一個(gè)全新的技術(shù)名詞更值得關(guān)注的是它背后那一整套真正解決數(shù)據(jù)工程問(wèn)題的能力。從 Spark 集群搭建、SQL 分析、外部數(shù)據(jù)源集成到 OOM 排查、信創(chuàng)數(shù)據(jù)庫(kù)適配這套技術(shù)棧在可預(yù)見(jiàn)的未來(lái)里依然會(huì)是大數(shù)據(jù)平臺(tái)的中堅(jiān)力量。對(duì)于正準(zhǔn)備入門 Spark 的讀者建議按本文順序先把 Standalone 集群搭起來(lái)跑通第一個(gè)spark-submit任務(wù)再逐漸嘗試用 Spark SQL 替代日常的 Hive 查詢。對(duì)于已經(jīng)在生產(chǎn)環(huán)境使用 Spark 的工程師可以對(duì)照第六節(jié)和第八節(jié)做一次資源參數(shù)體檢看看是不是還有 Task 傾斜、內(nèi)存浪費(fèi)的情況。技術(shù)熱榜來(lái)來(lái)去去但底層計(jì)算引擎的工程能力始終是數(shù)據(jù)團(tuán)隊(duì)最穩(wěn)的護(hù)城河。