开头
写这个系列的前几篇时,我都是以“把环境跑起来、把任务提交上去”为主线,到了第七篇,思路得换一换了。Scala 和 Spark 这套技术栈,真正难的不是 API 怎么调,而是你搭好集群、提交作业之后,那些藏在日志里、参数里、资源调度里的“暗坑”。这一篇就围绕我最近在实际项目里踩过的坑和调优记录来展开:从 Scala 安装与依赖拉取慢这种最磨人的小事,到 Spark 内存模型怎么配才合理,再到 Spark on YARN 上 CPU 只分配 1 个核这种经典到可以进面试题的疑难杂症,最后聊聊 Spark SQL 在外部数据源适配和真实数据分析场景里的玩法。
内容偏实战,不是入门教程。如果你已经能把 Spark 作业跑起来,但总觉得哪里不对劲、资源利用上不去、任务老是莫名其妙失败,那这篇就是写给你看的。刚接触 Scala 和 Spark 的新手也可以看,里面有不少环境搭建和配置的细节,能帮你少走弯路。
1. 环境搭建被忽略的细节,最消耗耐心
很多人觉得环境搭建是最简单的一步,其实恰恰相反。我在帮团队搭开发环境的时候发现,真正让人崩溃的往往不是 Spark 本身,而是 Scala 工具链和依赖下载这些“小事情”。
1.1 Scala 安装与 coursier 下载慢的破解思路
Scala 官方现在推荐用 coursier 来安装和管理 Scala 环境,这玩意儿本身设计得很好,能帮你管理多个 Scala 版本、快速切换,但它有个绕不开的问题——默认从 Maven 中央仓库拉取依赖,在国内网络环境下那速度简直感人。很多朋友卡在cs setup这一步半天没动静,误以为死机了。
这里分享一个我实测有效的做法:给 coursier 配置国内镜像源。在用户主目录下找到.config/coursier目录,编辑或新建config.properties文件,加上阿里云的 Maven 镜像地址。这样再执行cs setup,下载速度完全是两个世界。如果你是在公司内网环境,还可以用公司私有 Nexus 或 Artifactory 仓库,配置方式同理,把地址换成内网仓库就行。
还有个常见问题是 Scala 版本和 Spark 版本的匹配。Spark 3.x 系列分别支持 Scala 2.12 和 2.13,但同一个 Spark 发行包只能对应一个 Scala 主版本。你拿 Spark 3.3.0 的源码包去编译,默认用的是 Scala 2.12;如果你项目里用的是 Scala 2.13,那就要找对应的spark-3.3.0-bin-hadoop3-scala2.13这样的发行包。我见过不止一个同事在这里栽跟头,本地写代码好好的,一打包提交集群就报NoClassDefFoundError,最后发现就是 Scala 版本不一致导致的。
提示:在 pom.xml 或 build.sbt 里引入 Spark 依赖时,务必确认
scala.binary.version和集群上 Spark 发行包的 Scala 版本完全一致。
1.2 Spark 集群搭建的核心步骤与踩坑记录
搭建 Spark 集群本身并不复杂,Standalone 模式下一主两从就是改三个配置文件的事:spark-env.sh里配SPARK_MASTER_HOST、SPARK_WORKER_CORES、SPARK_WORKER_MEMORY,workers文件里写上从节点主机名,然后启动就行。但有几个细节很容易被忽略:
第一,主节点和从节点之间的 SSH 免密登录必须提前配好,否则sbin/start-all.sh会卡在密码输入上,甚至静默失败。第二,spark-env.sh里的JAVA_HOME要写绝对路径,不要依赖/etc/profile里的环境变量,因为 Spark 的启动脚本在某些发行版上不会加载登录 shell 的环境变量。第三,如果节点之间时钟不同步,Executor 的心跳上报会频繁出现超时,表现就是任务偶发性失败,排查起来非常隐蔽,建议集群规模稍大就部署一套 NTP 时间同步。
另外,如果是生产环境,我强烈建议直接走 Spark on YARN 的模式,不要再自己维护 Spark 集群的 Master/Worker 进程。YARN 作为统一的资源调度层,能把 Spark、MapReduce、Flink 这些计算框架的资源统一管起来,运维成本低得多。后面会专门讲 YARN 模式下 Spark 资源分配的一些坑。
2. Spark 内存模型,调优必须先懂原理
内存相关的参数是 Spark 配置里最容易被乱调的。很多人一上来就把spark.executor.memory调到很大,结果作业反而更慢,甚至频繁 Full GC。不把内存模型搞清楚,调参就是在瞎蒙。
2.1 统一内存管理机制详解
Spark 从 1.6 开始用统一内存管理(Unified Memory Management),把 Executor 的堆内内存分成了三块:Reserved Memory(预留内存)、User Memory(用户内存)、Spark Memory(Spark 内存)。其中 Spark Memory 又被 Execution Memory(执行内存)和 Storage Memory(存储内存)共享,两者之间可以互相抢占。
我用一个生活化的类比来帮你理解这个机制。想象你有一个双人书房,中间没有硬隔断,只有一块可以移动的屏风。Execution Memory 就像是你在书桌上铺开的资料,跑 Shuffle、Join、排序这些计算操作需要临时堆放数据;Storage Memory 就像是你的书柜,用来缓存 RDD、广播变量等数据。屏风的作用就是:书桌上放不下了,可以往书柜那边挪;书柜不够用了,也可以往书桌这边借。但有个规矩,书桌上的东西(执行内存)可以随时抢占书柜的空间,而书柜里的东西(存储内存)能不能抢占书桌,得看书桌够不够宽松——如果执行内存本身还很宽裕,书柜可以借一点;如果执行内存已经吃紧了,书柜就只能被挤。
这套机制的好处是:Shuffle 等计算密集型操作对内存的需求往往是突发的、刚性的,内存不够直接就会溢写到磁盘,性能断崖式下跌;而缓存类数据是弹性的,被挤掉大不了后面重新计算一次。所以设计上让执行内存拥有更高优先级。
2.2 关键参数与内存配置实战
具体到参数上,核心就三个:spark.memory.fraction、spark.memory.storageFraction、spark.executor.memory。spark.memory.fraction默认是 0.6,表示 Spark Memory 占整个堆内存的比例,剩下 0.4 留给 User Memory 和 Reserved Memory。spark.memory.storageFraction默认是 0.5,表示 Storage Memory 占 Spark Memory 的初始比例,另一半给 Execution Memory。
我见过很多人把spark.memory.fraction调到 0.8 甚至 0.9,觉得这样 Spark 能用的内存更多了。这个思路在纯缓存场景下说得通,但如果你同时有大量的 Shuffle 操作,User Memory 被压得太小会导致序列化、对象存储等操作频繁 GC。我个人的经验值是:混合负载场景保持默认 0.6 就行,真要调也不要超过 0.75。
举个实际的配置例子。假设一个 Executor 分配了 4GB 堆内存(spark.executor.memory=4g),默认参数下:
- Reserved Memory 固定 300MB,不太受关注
- User Memory 大约 (4GB - 300MB) × 0.4 ≈ 1.48GB,用来存用户数据结构、防止 OOM 的缓冲
- Spark Memory 大约 (4GB - 300MB) × 0.6 ≈ 2.22GB
- 其中 Storage Memory 初始约 1.11GB,Execution Memory 初始约 1.11GB
如果你的作业以 Cache 和广播变量为主,可以把spark.memory.storageFraction调到 0.6 到 0.7;如果是 TPC-DS 这类复杂 SQL 查询,Shuffle 满天飞,我建议把spark.sql.shuffle.partitions调小一些来控制输出文件数,而不是去动storageFraction。
还有一个容易忽略的点:堆外内存。spark.memory.offHeap.enabled默认是 false,但如果开了spark.memory.offHeap.size,这部分内存也会计入统一内存管理,和堆内内存共享同一套memory.fraction分配逻辑。在 YARN 模式下,堆外内存还会受到spark.executor.memoryOverhead的影响,这个参数默认是 executor 内存的 10%,最小 384MB,用来容纳 JVM 自身开销、线程栈、NIO Buffer 等。如果你的作业用到了大量的 Netty 通信或直接内存,memoryOverhead不够会报Container killed by YARN for exceeding memory limits,这时候要果断调大,要调大到多少,看实际情况定,我一般从 1GB 起步试。
3. Spark on YARN 的 CPU 配置疑难杂症
热搜词里有这么一条:“spark on yarn cpu只能用1个是为什么”。看到这个词条我第一反应就是太真实了,这个问题我在不同场合被问了不下十次。
3.1 问题现象与根因定位
现象一般是这样:在 YARN 上提交 Spark 作业,打开 Spark Web UI,发现每个 Executor 只有 1 个 vCore,任务并行度上不去,整个作业慢得像蜗牛。很多人第一反应是去调spark.executor.cores,改成 4,重启作业,结果发现还是 1 个核。
根因十有八九出在 YARN 的调度器配置上。YARN 默认的容量调度器(Capacity Scheduler)里,有一个参数叫yarn.scheduler.capacity.maximum-am-resource-percent,控制的是 ApplicationMaster 最多能拿到集群资源的比例,但和这个 CPU 问题直接相关的是另外一类限制。如果你用的是 Fair Scheduler,还要看yarn.scheduler.fair.maximum-am-resource-percent。但最最常见的,其实是yarn.nodemanager.resource.cpu-vcores这个参数没有配置或者配置偏低。
举个例子。很多机器是 16 核的,但你安装好 Hadoop 后没有去改yarn-site.xml,导致 NodeManager 上报给 ResourceManager 的 vCore 数量是默认的 8,甚至在某些发行版里可能更少。然后你在 Spark 里申请 4 个 Executor、每个 4 核,结果 YARN 一算,总共只有 8 个 vCore 可用,你的 ApplicationMaster 还要占 1 个,剩下 7 个,你每个 Executor 要 4 个核,那只能起来 1 个 4 核 Executor,剩下 3 个 Executor 等不到资源就不断重试。如果你设置了spark.executor.cores=1,那你最多就能跑 7 个 Executor,每个 1 核,看起来就是“怎么只有 1 个核”。
3.2 排查步骤与解决方案
遇到这个问题,我建议按这个顺序排查:
- 查看 YARN 资源监控页面,确认集群总的 vCore 数。如果不到物理核数的一半,基本就是
yarn.nodemanager.resource.cpu-vcores没配好。 - 检查调度器的最大 AM 资源比例。如果你跑了多个作业共享队列,AM 占用的资源比例超过了上限,作业会一直处于 ACCEPTED 状态,请求不到容器。
- 检查 Spark 任务提交时设置的
--executor-cores和--total-executor-cores。前者决定单个 Executor 的核数,后者控制整个作业的总核数上限。 - 最后确认
spark.task.cpus的配置。这个参数默认是 1,表示每个任务占用 1 个核。如果它被改成大于 1 的值,比如 2,那么一个 4 核的 Executor 里最多只能同时跑 2 个任务,并发度直接减半。
我再补充一个容易踩坑的点:spark.dynamicAllocation.enabled开启时,Spark 会根据任务负载动态调整 Executor 数量。这个功能在 YARN 上用的时候,必须同时开启spark.dynamicAllocation.shuffleTracking.enabled,否则在 Shuffle 密集的作业中,动态缩容会把承载 Shuffle 数据的 Executor 给回收掉,导致下游任务重新拉取数据,代价很大。
注意:修改
yarn-site.xml中的 CPU 和内存配置后,必须重启 NodeManager 才会生效。别问我是怎么知道的,我改完参数忘了重启,排查了整整一个下午。
下面给出一个我在生产环境验证过的基础配置参考表。
| 配置项 | 推荐值 | 说明 |
|---|---|---|
yarn.nodemanager.resource.cpu-vcores | 物理核数 × 1~1.5 | 如果开了超线程且有其他业务混部,取物理核数即可 |
yarn.nodemanager.resource.memory-mb | 物理内存 × 0.8 | 预留一些给系统本身和其他进程 |
spark.executor.cores | 3~5 | 太小并行度不足,太大容易引发 GC 抖动,4 是个比较稳妥的值 |
spark.executor.memory | 受单节点总资源和 Executor 数量约束 | Executor 内存 + overhead 不要超过 NodeManager 剩余内存 |
spark.task.cpus | 1 | 除非任务里有耗 CPU 的复杂计算,否则保持默认即可 |
4. Spark 日志与默认配置的实战解读
每份 Spark 日志里都会出现一行Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties。这个信息看着像警告,其实只是 Spark 在告诉你它没找到你的自定义日志配置,正在用默认配置。但很多人被它误导,以为系统出问题了。
4.1 log4j 配置与日志分级策略
Spark 3.x 用的是 log4j 2.x,默认配置会以 INFO 级别输出日志。对于生产环境来说,INFO 级别的日志量在 Shuffle 和任务频繁调度时会非常可观,尤其是org.apache.spark.scheduler和org.apache.spark.executor这两个包,几乎每个任务启动和结束都会打日志。如果你的 Spark History Server 长期运行,日志文件会膨胀得非常厉害。
我一般在提交作业时通过--conf spark.driver.extraJavaOptions=-Dlog4j.configurationFile=file:///opt/conf/log4j2.xml来指定自定义配置文件,或者把log4j2.xml直接放到 Spark 的conf目录下。配置的核心思路是分级别定义 Appender:WARN 级别的写入滚动文件,专门用来排查问题;INFO 级别的输出到控制台,方便开发调试;对特定第三方包如org.apache.hadoop可以单独调高到 WARN,减少无关刷屏。
关于日志里最常见的报错,我可以列一个速查表:
| 日志关键字 | 通常原因 | 解决方案 |
|---|---|---|
Container killed by YARN for exceeding memory limits | Executor 总内存超过容器上限 | 增大spark.executor.memoryOverhead |
Lost executor+FetchFailedException | 网络抖动、BlockManager 通信失败或磁盘故障 | 检查节点网络与磁盘,开启spark.shuffle.service.enabled |
java.lang.OutOfMemoryError: Java heap space | 堆内存不够 | 增大spark.executor.memory,优化数据分区大小 |
Serialized task exceeded max allowed size | 任务闭包过大,常因 Driver 端把大对象传给了 Executor | 使用广播变量,精简闭包 |
Initial job has not accepted any resources | 资源不足,作业等不到容器 | 检查队列资源和 AM 比例配置 |
4.2 日志驱动的问题排查方法论
日志不仅是记录,更是排查问题的入口。这里分享一个我自己长期使用的排查套路:当一个 Spark 作业失败后,不要先看堆栈,而是从上到下按时间线顺序捋日志。第一步看 Driver 日志里的 ERROR 和 WARN,第二步在 Executor 日志里搜TaskSetManager报的Lost task信息,第三步去 YARN 的 ResourceManager 日志看容器被 kill 的具体原因。大多数疑难杂症都能通过这三步定位到方向。
举一个真实案例。有一次我负责的离线任务频繁失败,错误信息很杂,有时是FetchFailedException,有时是Connection refused。起初以为是网络故障,排查了一圈发现节点网络正常。后来看了 NodeManager 的日志才发现,执行节点的磁盘空间被写满了。原因是我的作业里有一个大的.cache()操作,缓存数据落盘到/tmp,而/tmp挂载的分区只有 20GB。问题的根因不是网络,而是磁盘空间不足导致 BlockManager 无法写入临时数据。这种问题的隐蔽之处在于,Spark 的报错信息不会直接告诉你“磁盘满了”,只会表现为各种网络和 IO 异常,只有靠日志逐层排查才能找到根因。
5. Spark SQL 在真实业务场景中的玩法与优化
Spark SQL 是日常开发中使用频率最高的模块之一。前面几讲我多少提到过 DataFrame API,这一讲来深入聊聊 Spark SQL 在外部数据源适配和实际业务案例分析中的表现。
5.1 外部数据源适配:以达梦数据库为例
国内不少政企项目会用达梦数据库(DM),Spark 官方并没有提供达梦的 JDBC 连接器,但这难不倒我们。JDBC 数据源天然就是这样用的:通过spark.read.format("jdbc")配上达梦的驱动类dm.jdbc.driver.DmDriver和连接 URLjdbc:dm://host:5236,就能让 Spark SQL 像读 MySQL 一样读达梦的数据。
不过有几个细节必须处理好。第一,驱动 jar 包要同时分发到每个 Executor 节点,最省事的方式是提交作业时用--jars参数带上,或者在spark-env.sh里配置SPARK_CLASSPATH。第二,达梦的 SQL 方言和 MySQL 有些差异,Spark SQL 做谓词下推时会生成 JDBC 查询语句,有些函数可能达梦不支持,会导致下推失败甚至报错。稳妥的做法是:对达梦表不要直接做复杂函数运算,先通过filter下推基本条件,把数据拉到 Spark 再处理,性能虽然差一点,但兼容性最稳。第三,注意配置partitionColumn和numPartitions来并行读取大表,否则默认单分区读取,数据量大时效率极低。
我实测过一张千万级别的表,通过partitionColumn按自增 ID 分为 8 个分区读取,耗时从单分区的 40 多分钟降到了 7 分钟左右。代价是每一个分区都会建立一个 JDBC 连接,达梦数据库的连接数上限要提前调好,不然会把数据库连接池打爆。
5.2 Spark SQL 优化的核心原则:先看计划
说到 Spark SQL 优化,我养成的一个习惯是:任何慢查询来了,第一件事不是调参数,而是看一眼执行计划。用df.explain(true)能看到逻辑计划和物理计划,重点观察三件事:过滤条件下推到了哪个阶段、有没有不必要的 Shuffle、Join 用的是 SortMergeJoin 还是 BroadcastJoin。
举一个典型的优化案例。一张订单表和一张用户表做 Join,订单表有几亿行,用户表只有几千行。默认情况下 Spark 会走 SortMergeJoin,两边都要排序和 Shuffle,耗时非常难看。但只要提前把用户表用小表广播出去,spark.sql.autoBroadcastJoinThreshold调大到合适值,或者在建表时用hint指定 broadcast join,执行计划立刻变成无需 Shuffle 的 BroadcastHashJoin,性能提升一个量级。
还有一个经常被忽视的参数是spark.sql.shuffle.partitions,默认 200。如果你的数据量不大,200 个 Shuffle 分区会把每个分区塞得很小,反而增加调度开销。反之,如果你的数据量很大,200 个分区又会导致每个分区上百 GB,单任务内存压力巨大。我的做法是结合spark.sql.adaptive.enabled一起用,开启 AQE 之后 Spark 可以在运行时自动合并小分区、自动调整 Join 策略,配合合理的初始shuffle.partitions,能达到相当好的效果。
5.3 跨界案例分析:实车试验大数据与新能源能量管理
热搜词里有一条和汽车行业强相关的词条:“基于实车试验大数据分析的插电式混合动力汽车能量管理策略解析”。这个案例虽然看起来离互联网技术很远,但底子和 Spark 数据分析的套路完全一致。插电式混合动力汽车在运行时,整车控制器(HCU)会在发动机和电机之间做能量分配,这一过程会产生海量的时间序列数据:车速、加速度、电池 SOC、发动机转速、扭矩、温度等。对这些数据分析得越深入,能量管理策略就能做得越精细。
有一个很典型的分析思路:先按时间窗口对实车试验数据做切分,比如以 10 秒为一个片段,然后用聚类算法把这些片段划分成不同的工况类型,比如城市拥堵、郊区工况、高速巡航。再针对每一类工况,统计发动机和电机的能量分配比例、电池充放电规律,找出效率最低的区间,结合能量管理策略需求去优化控制参数。这里面的数据清洗和特征工程环节,正是 Spark SQL 和 Spark MLlib 的强项。比如在数据预处理阶段,用窗口函数对车速做滑动平均平滑掉传感器噪声,用when条件判断把 SOC 越界的异常点剔除,这些操作在 Spark SQL 里都只是几个表达式的事。
这个案例给我们的启发是:Spark 大数据分析的应用边界远超互联网广告推荐,也不止是日志分析,任何能产生海量数据的行业——汽车、能源、制造、金融——都可以用同一套工具去挖掘规律。
6. 高频面试题与工程经验速查
很多朋友面试时会背一些 Spark 原理题,比如 RDD 和 DataFrame 的区别、宽依赖和窄依赖、Stage 划分机制。这些概念性的东西当然重要,但近几年面试官越来越喜欢问场景题,这里我也整理几个我认为最有代表性的问题,结合前面的实战内容做一个串联式的复盘。
6.1 为什么我的执行计划会出现大量 Shuffle
这个问题要回答清楚,需要理解窄依赖和宽依赖的本质区别。窄依赖指父 RDD 的每个分区最多被子 RDD 的一个分区使用,而宽依赖指一个分区的数据会被打散到多个子分区,这就必然产生 Shuffle。以 GroupByKey 和 ReduceByKey 为例,前者会把 V 原样 shuffle 到下游再合并,后者在 map 端就先做一次预聚合,减少网络传输量,这就是为什么 spark 官方和所有面试答案都推荐用 ReduceByKey 的原因。
但实际工程里,我在代码 review 时发现很多人并不是不懂这个理论,而是用 DataFrame API 时根本不关心底层用了什么算子。DF 的执行计划是 Catalyst 优化器生成的,它会在你写完代码后自动应用谓词下推、列剪裁、常量折叠等优化手段。所以如果你想减少无关 Shuffle,与其纠结算子选型,不如先确保数据模型设计合理:分区键是否均匀、文件大小是否合理、过滤条件是否能下推。
6.2 Executor、Core、Task、并行度之间的关系
这是面试里最容易被问晕的一组概念。我习惯用一个类比来讲:Executor 相当于一台独立的“工人宿舍”,它独占一块内存和若干核;Core 相当于宿舍里的“床铺”,有多少张床决定了这个宿舍能同时住几个人;Task 相当于“工人”,每个工人需要一张床才能开工;并行度就是整个集群里所有宿舍床铺数量的总和,直接影响同一时间能有多少任务并发执行。
在 YARN 模式下,每个 Executor 是一个 Container,能用的核数由--executor-cores决定,能用的内存由--executor-memory决定。Container 能够申请的核数和内存上限,受制于 NodeManager 的资源量和调度器的配置。很多时候作业跑得慢,不是没资源,而是资源隔离和配置没匹配上。理解了这层关系,你再回看第三部分的 CPU 核数问题,逻辑就非常清晰了。
6.3 内存调优时最容易被忽略的隐性开销
前面讲内存模型时提到spark.executor.memoryOverhead,这里再展开说一个隐性开销:在 Spark SQL 场景下,Java 对象在堆里占的空间远大于你以为的大小。一个字符串在 JVM 里除了 char 数组,还有对象头、长度字段、对齐填充这些额外开销,整体占用往往是数据本身大小的 2 到 3 倍。这就是为什么用 Kryo 序列化能大幅降低内存占用,因为它把对象变成了紧凑的字节数组。
所以如果你发现 Executor 的实际内存利用率很低,但 GC 却很频繁,可以考虑把spark.serializer改成org.apache.spark.serializer.KryoSerializer,并注册需要的类。另外,处理大表时优先用 DataFrame API 而不是 RDD API,因为 DataFrame 的内存布局基于列式存储,对内存的利用率远高于普通 Java 对象。
最后再分享一个我个人的心得:以前做 Spark 调优的时候,总喜欢把所有参数都审视一遍,逐个去调。后来经历得多了,发现大多数问题集中在这几件事上:数据倾斜、Shuffle 量过大、资源没配够、序列化开销大。先把执行计划和数据规模摸透,再动手调参数,往往能一针见血。遇到奇怪的资源分配问题,记得从 YARN 的调度配置查起,而不是一头扎进 Spark 参数里反复试。做大数据分析,花在数据理解上的时间,从来都不会白费。