ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Spark on YARN 实战调优:从内存模型到 CPU 核数疑难杂症

Spark on YARN 实战调优:从内存模型到 CPU 核数疑难杂症 开头写这个系列的前几篇时我都是以“把环境跑起来、把任务提交上去”为主线到了第七篇思路得换一换了。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_MEMORYworkers文件里写上从节点主机名然后启动就行。但有几个细节很容易被忽略第一主节点和从节点之间的 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 MemorySpark 内存。其中 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.memory4g默认参数下Reserved Memory 固定 300MB不太受关注User Memory 大约 (4GB - 300MB) × 0.4 ≈ 1.48GB用来存用户数据结构、防止 OOM 的缓冲Spark Memory 大约 (4GB - 300MB) × 0.6 ≈ 2.22GB其中 Storage Memory 初始约 1.11GBExecution 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.cores1那你最多就能跑 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.cores3~5太小并行度不足太大容易引发 GC 抖动4 是个比较稳妥的值spark.executor.memory受单节点总资源和 Executor 数量约束Executor 内存 overhead 不要超过 NodeManager 剩余内存spark.task.cpus1除非任务里有耗 CPU 的复杂计算否则保持默认即可4. Spark 日志与默认配置的实战解读每份 Spark 日志里都会出现一行Using Sparks 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.configurationFilefile:///opt/conf/log4j2.xml来指定自定义配置文件或者把log4j2.xml直接放到 Spark 的conf目录下。配置的核心思路是分级别定义 AppenderWARN 级别的写入滚动文件专门用来排查问题INFO 级别的输出到控制台方便开发调试对特定第三方包如org.apache.hadoop可以单独调高到 WARN减少无关刷屏。关于日志里最常见的报错我可以列一个速查表日志关键字通常原因解决方案Container killed by YARN for exceeding memory limitsExecutor 总内存超过容器上限增大spark.executor.memoryOverheadLost executorFetchFailedException网络抖动、BlockManager 通信失败或磁盘故障检查节点网络与磁盘开启spark.shuffle.service.enabledjava.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 外部数据源适配以达梦数据库为例国内不少政企项目会用达梦数据库DMSpark 官方并没有提供达梦的 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 参数里反复试。做大数据分析花在数据理解上的时间从来都不会白费。
返回列表