ARTICLE DETAIL

资讯详情

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

Spark性能调优实战:从资源调度到数据倾斜的完整优化指南

Spark性能调优实战:从资源调度到数据倾斜的完整优化指南 1. 先别急着调参部署形态和资源分配决定调优上限凌晨两点的报警群里突然弹出一条消息数仓的夜批作业又超点了一个Stage卡了快四十分钟没动静。这种场景做过大数据的都懂——代码跑了大半年没动过数据量一上来Task耗时开始分叉从几十秒拉到十几分钟整个调度链跟着雪崩。我当时的第一反应不是急着堆executor内存而是把Spark UI打开看调度器在等什么、Stage内部的时间分布有没有长尾、Shuffle Read的数据量是不是突然暴涨。正是这一套从现象到参数的排查习惯让我慢慢把一批固定作业从5小时跑到了1小时以内。这篇文章我会把Spark调优的实操经验拆开讲部署形态选择、内存模型、数据倾斜、Spark SQL运行时优化、监控与GC、以及一套可复制的调优流程。适合正在做离线数仓、实时链路批处理、或者被作业慢、Executors被OOM杀死、数据倾斜拖垮Stage折磨的工程师。内容偏生产实践不堆概念每一步都给出为什么这么调的理由。1.1 YARN、Standalone该选谁资源调度的现实约束生产环境里我绝大多数情况会选择YARN作为Spark的资源管理器而不是Standalone。原因很简单Standalone是Spark自己管资源任务和集群里其他计算框架之间没有统一调度天然没法和其他团队共享物理机而YARN把CPU和内存抽象成Container配合Hadoop生态可以做到多租户隔离、队列配额、优先级调度。你在Standalone上调出来的参数换到YARN上可能完全失效因为资源申请方式不同。另一个容易忽略的约束是机架感知。Spark在YARN上跑时Executor分配会优先考虑输入数据所在的节点减少跨网络拉数据的代价。这个在Standalone里也能做但配置成本和稳定性不如YARN自带的那套。所以我的判断标准是只要公司里已经有HDFS、Hive等组件老老实实走YARN只有小规模测试环境或者纯业务验证机才考虑Standalone。1.2 动态资源分配是必选项前提是开启Shuffle Service很多人部署Spark集群时习惯一次性指定executor数量比如固定100个Executor每个8G。这个做法在业务低峰期纯属浪费高峰期又容易被其他任务挤掉。我的建议是优先开启动态资源分配让集群按作业实际需求伸缩Executor规模。spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 2 \ --conf spark.dynamicAllocation.enabledtrue \ --conf spark.dynamicAllocation.initialExecutors5 \ --conf spark.dynamicAllocation.minExecutors5 \ --conf spark.dynamicAllocation.maxExecutors50 \ --conf spark.shuffle.service.enabledtrue注意最后一行必须开启Shuffle Service。这是动态分配能安全回收Executor的关键。Shuffle数据在Spark架构里是写在本地磁盘的如果某个Executor被回收下线它的Shuffle中间结果可能直接丢掉下游Task再想读就得失败。开启Shuffle Service后数据由统一的常驻服务管理Executor就算被回收也不影响Shuffle完整性。我见过不少团队直接抄网上的动态分配配置但漏了shuffle service结果任务跑到一半反复报FetchFailed最后只能关掉动态分配。1.3 一套能直接上生产的初始资源参数说到具体参数我一般不会让Executor内存超过单机物理机的1/2。比如物理机64GExecutor内存给16G一台机器最多放3个Executor。如果给到32G一旦JVM堆内GC做Full GC停顿时间会明显拉长还会影响同机器的其他容器。下面这套是我日常起新作业时的基准配置你先照这个跑再根据现象做微调参数建议值说明spark.driver.memory2g-4gDriver一般不承担重计算但如果有collect操作会吃内存需上调spark.executor.memory8g-16g建议根据单条记录的复杂度和清洗逻辑调整spark.executor.cores2-3不宜过大单个Executor并发太多Task会加剧GC和线程切换spark.sql.shuffle.partitions200起调根据shuffle数据量动态修正后面会详细讲spark.default.parallelism与总核数匹配建议设为executor数×cores数spark.dynamicAllocation.enabledtrue多租户集群必须开配合shuffle service这套基准值的原则是保守起步避免一上来就把资源顶满。调优调优先得有稳定的基线再找瓶颈。你连基线都没跑顺后面的内存、倾斜、SQL优化都无从对比。2. Spark 3.x内存模型拆解MemoryOverhead才是隐藏的OOM推手很多刚开始做Spark调优的人有个误区一看到Executors OOM就疯狂增加spark.executor.memory结果内存调到20G、24G任务照样被杀。他们没意识到YARN判定一个Container是否超内存看的是整个进程的物理内存占用而不仅仅是你配置的JVM堆。这里面藏着一个关键角色——MemoryOverhead。2.1 堆内分配不是全部execution、storage、reserved怎么分Spark 3.x的统一内存管理模型里一个Executor的JVM堆内存被分成三块Execution内存负责Shuffle、Join、Aggregation等临时数据Storage内存负责缓存RDD、Broadcast数据Reserved内存是系统保留默认300MB用来跑Spark内部逻辑。// 统一内存管理核心参数 spark.memory.fraction 0.6 // Execution Storage可用堆的比例 spark.memory.storageFraction 0.5 // Storage在共享区域内占的比例也就是说假如Executor堆为16G实际可用于执行和缓存的内存只有9.6G左右(16G * 0.6)其中Execution和Storage各有4.8G的默认配额。但这两个区域可以互相借用Execution用满了可以挤占Storage的空间反之亦然。理解这个机制很重要比如你如果手动cache了大量DataFrame又跑一个需要大内存的Join两者会互相争夺空间极端情况触发频繁Spill甚至OOM。2.2 被YARN杀掉的Executors从日志里认出真正的堆外OOM这是我踩过最深的一个坑。某次清洗任务跑在大批量JSON解析上Executor日志里几乎没有堆内OOM报错GC时间也正常但YARN反复报告Container killed for exceeding memory limits。当时第一反应是怀疑spark.executor.memory不够从8G一路加到24G问题依旧。后来查Container的物理内存监控曲线发现JVM堆只用到了一半但进程RSS已经顶满上限——问题出在堆外。YARN会在你申请的JVM堆之外额外追加一块MemoryOverhead内存默认参数是max(384MB, 0.1 * spark.executor.memory)。如果你的作业里用到大量堆外资源比如读取Parquet/ORC时的Native内存、某些C扩展库、或者Netty的堆外Buffer这部分内存会从Overhead里扣。JSON解析场景特别典型因为解析过程中频繁产生临时Native内存和字符串对象堆内看似够用堆外早就爆了。碰到这种场景先加Overhead内存spark-submit \ --executor-memory 8g \ --conf spark.yarn.executor.memoryOverhead1g注意这里数值是绝对值不是比例。从16G堆调到24G堆解决不了堆外问题但把Overhead从默认0.8G提升到1.5G问题反而立刻消失。这个案例后来我写进团队的排障手册里凡是用到自定义序列化、大量JSON解析、或者堆外缓存的项目MemoryOverhead起步就按1G算而不是信默认0.1倍。2.3 Kryo序列化与内存占用减少堆内压力的低成本方案默认情况下Spark使用Java序列化对象序列化之后体积大反序列化时临时对象也多这套机制在内存紧张的业务里非常吃亏。我处理大状态作业时基本都会切到Kryo序列化配合类注册减少元数据开销。--conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.kryo.registrationRequiredtrue \ --conf spark.kryo.registratorcom.example.MyKryoRegistratorKryo的性能优势很明显序列化速度快占用空间大约是Java序列化的一半到三分之一。代价是需要自己注册类尤其是算子闭包里涉及的自定义对象、case class漏了注册直接报错。我建议把Kryo registrationRequired设为true宁可启动时发现漏注册而快速暴露问题也不要等运行到深层算子时突然抛异常。这个改动对纯Shuffle任务提升不大但对Broadcast变量大、或者反复Checkpoint的作业收益非常明显。3. 数据倾斜和Shuffle一次拖垮整个Stage的经验复盘数据倾斜在Spark作业里的典型症状是一个Stage只有几十个Task其他Task几十秒跑完某一两个Task要跑十几分钟而且Shuffle Read的数据量比平均值高出好几个量级。这种问题靠调内存基本没用必须从数据和代码层面去拆解。3.1 定位倾斜Stage先学会看Spark UI的时间分布打开Spark UI看Stages页优先关注两列Shuffle Read Size/Records 和 任务执行时间分布。正常情况的Task执行时间应该大致集中如果出现明显的长尾且长尾Task对应的Shuffle Read Size异常大基本可以确定是倾斜。我复盘过的典型case是订单表按城市分组后做topN统计一个一线城市的key占比超过总量40%导致单个Task处理的记录数远超其他Task。这时候executor加核、加内存都没用因为瓶颈在单Task的处理能力上而不是整体吞吐量。先定位到具体Stage和具体Key再对症下药后面的策略才能真正生效。3.2 groupBy型倾斜两步加盐聚合针对groupBy/聚合导致的倾斜最实用的方案是加盐两阶段聚合。思路是先给倾斜的Key加一个随机前缀把一个大Key拆成多个小Key并行聚合再在下一层去掉前缀做全局聚合。// 阶段一对倾斜Key加随机前缀 val salted df .withColumn(salt, when($key.isin(tiltedKeys: _*), rand() * 10).otherwise(lit(0))) .withColumn(saltedKey, concat($salt, lit(_), $key)) .groupBy(saltedKey) .agg(sum(amount).as(partAmount)) // 阶段二去掉前缀全局再聚合一次 val result salted .withColumn(realKey, split($saltedKey, _).getItem(1)) .groupBy(realKey) .agg(sum(partAmount).as(finalAmount))这里有两个前提你要注意。一是你得先确定哪些Key是倾斜的可以通过SQL统计各Key的量级而不是对全部Key都加盐二是加盐粒度不要过大8到16路拆分足够拆太细反而增加第二阶段聚合的压力。我在线上实测中加盐后那个拖后腿的Task耗时从18分钟下降到2分钟左右整体Stage缩减了70%以上。3.3 join型倾斜广播Join和三张前置清单Join型倾斜和GroupBy型不太一样它不是单Key数据大而是两个表关联时某个Key在右侧表里成百倍放大。最简单的解法是广播小表让每个Executor本地持有完整的小表数据完全绕开Shuffle。SELECT /* BROADCAST(dim) */ f.order_id, d.user_level FROM fact_table f LEFT JOIN dim_table d ON f.user_id d.user_id但要小心广播Join并非没有成本。我把这个策略归纳成三张前置清单一小表数据量是否在广播阈值范围内默认是10MB超过的话可以临时调大但调太大Driver端会承受巨大压力二小表更新频率如何如果每五分钟变一次广播任务频繁刷新反而拖慢整个作业三小表是否做了列裁剪我见过不少人直接把整张宽表广播出去动辄上百MB后端网络和GC直接被打爆。如果大表join大表单边倾斜用AQE的倾斜处理机制会顺手很多这在第4章展开。手动方案通常是先过滤无效Key、再对倾斜Key单独Join、最后做Union合并但这套代码维护成本高建议能交给引擎处理的就不手写。3.4 分区数不是越大越好写文件和Shuffle的平衡最后一个高频问题分区数怎么设置。很多人误以为分区数越多并行度越高就把spark.sql.shuffle.partitions调到1000结果单个分区数据过碎写HDFS时小文件爆炸NameNode压力剧增任务不但没变快反而更慢。我的经验公式是先按Shuffle数据量估算每个分区目标大小100MB到200MB之间比较合适。假设Stage Shuffle写入量是50GB那50GB / 150MB约等于340个分区。你可以用这个作为spark.sql.shuffle.partitions的初始值再根据Task执行时间微调。还要留意repartition和coalesce的区别repartition会全量shufflecoalesce不做数据移动但只适合在分区数调小的场景。如果你只是想把1000个分区合并成50个用coalesce就够了如果是想把50个扩大到200个提高并行度那就只能repartition。4. Spark SQL的运行期干预AQE、CBO与代价模型Spark 3.x之后调优有了新的维度你不再只用代码去控制实际执行而是可以通过AQE和CBO让引擎在运行过程中自己调整执行计划。这一章讲怎么打开这些开关、什么时候它们会帮上大忙、什么时候反而帮倒忙。4.1 AQE为什么会自己改执行计划开开关之前的理解AQE的全称是Adaptive Query Execution它在每个Shuffle Stage完成后基于真实的输出数据量去优化后续的执行计划。开启后引擎能做三件重要事情动态合并小分区、动态调整Join策略、自动处理Join倾斜。--conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.sql.adaptive.advisoryPartitionSizeInBytes128MB \ --conf spark.sql.adaptive.skewJoin.enabledtrue这里的关键认知是AQE不是在编译期静态决定计划而是把决策延迟到运行期。比如某个大表和小表Join静态规划时因为小表的预估大小略大于广播阈值而选择了SortMergeJoin导致几十GB的Shuffle。有了AQE它会在Shuffle Map阶段完成后看到小表的真实大小只有8MB自动切换成BroadcastHashJoin那次Shuffle直接被干掉。建议在Spark 3.2及以上版本直接默认开启AQE几乎零成本。唯一需要注意的坑是动态合并分区时如果advisoryPartitionSizeInBytes设置太大会把本应并行执行的Task合并得过少导致单Task负担过重。我一般建议从128MB起步跑完看Task耗时再调整。4.2 广播Join阈值和手动Hint小表不小时怎么办AQE虽然能自动判断广播但静态规划阶段仍有一个阈值兜底spark.sql.autoBroadcastJoinThreshold默认10MB。这个参数对付看起来稍大的表仍然有效但我不建议盲目调大因为广播表要在Driver端收集再分发给所有Executor表越大Driver越容易被压垮。我线上通常把阈值控制在20MB以内。如果某张维度表确实有100MB但逻辑上必须全量参与Join可以针对单条SQL手动加Hint而不是全局调阈值。这等于明确了这张表这次就是要广播不会影响其他作业的默认行为。SELECT /* BROADCASTJOIN(dim) */ ... FROM fact f JOIN dim d ON f.dim_id d.id这里有一个很多人忽略的坑广播表在Executor端是驻留在Storage内存里的反复缓存和释放会直接影响GC。如果并发作业很多建议单独评估这个表能否用缓存表或宽表改造替代而不是每次运行都广播一遍。4.3 谓词下推、列裁剪从物理计划里揪出整表扫描代码层面调优最见效的是谓词下推和列裁剪。我记得一次排查同学写了一个DataFrame API任务把整张宽表读了进来再用filter做过滤结果每次跑30多分钟。我打开物理计划发现Filter在读文件之后才做列裁剪也没生效底层把整表的200多列全扫了一遍。优化后的做法是先做列裁剪只select需要的列过滤条件尽可能写在读取阶段。对于Parquet这类列式存储裁剪效果尤其明显IO量能降一个数量级。此外如果表按时间分区直接用partition过滤不要全表扫完再where。val result spark.read.parquet(/path/table) .filter($dt 2025-06-01) // 分区裁剪 .select(user_id, order_amount) // 列裁剪 .groupBy(user_id) .agg(sum(order_amount))这类代码改造无需魔法只是让引擎从源头上少读数据。很多线上任务慢在第一天写代码时贪了方便select全列、scan全分区后面加硬件参数是治标不治本。4.4 一个真实案例把2小时大查询拆成三段预聚合有一次业务侧给我一个复杂SQL涉及五张大表的Join和多次GroupBy跑一次接近2小时。打开执行计划发现中间有一个宽依赖非常重Shuffle数据量接近8TB。这种任务靠加参数解决不了唯一的路子是把计算拆成预聚合层。我把它拆成三段先对适合过滤的维度提前聚合缩小事实表关联前的规格再按核心维度做一次中间层聚合把明细数据压到1/10以下最后业务端点查直接读预聚合表。改造后同样的查询从2小时降到8分钟。这个案例给我的体会是SQL调优首先是数据模型和查询设计的调优其次才是Spark参数的调优。一个写得好的查询用20个Executor就能跑赢写得差的查询开100个Executor。5. 从UI数据到GC日志把监控变成调优证据调优最忌讳拍脑袋改一个参数之前必须先有数据支撑。Spark UI、事件日志、GC日志这三样东西是每个Spark作业排障时最直接的证据来源。5.1 每个排障者该固定的三个Spark UI入口Spark UI里面信息很多但真正高频使用就三个地方一是Jobs页里的Scheduling Delay。如果一个Task长期处于Running但迟迟不开始看调度等待时间。二是Stages页的Shuffle Read Size这是判断数据倾斜和Shuffle压力的主要指标。三是Executors页的内存和GC时间曲线。Task执行时间分布我会专门看有没有长尾GC时间占比超过总执行时间15%那就该去处理JVM了。现在很多集群通过独立History Server查看历史作业不只是看正在跑的任务。建议把事件日志打开作业结束后也能回放当时的执行计划、各Stage耗时、Task分布排障效率会高很多。spark.eventLog.enabledtrue spark.eventLog.dirhdfs://nameservice/spark-logs spark.history.fs.logDirectoryhdfs://nameservice/spark-logs5.2 事件日志和Dropwizard运维期间怎么回放现场除了UI的历史页面还可以利用Spark Metrics系统把指标推到外部的PrometheusGrafana监控看板。通过spark.metrics.conf配置可以把每个Executor的Shuffle读写量、GC次数、JVM堆使用率暴露给监控系统这样任务运行期间就能实时看到内存曲线。这在定位作业是一点点变慢还是突然掉坑时特别有用。我遇到过一种诡异问题任务每天新增一点数据某天触发了Shuffle数据量过大的临界点原本2G的Executor堆突然频繁GC。通过监控曲线对比能清楚看到堆内存从稳定60%突然冲到顶才知道要调整Overhead和Shuffle参数而不是盲目加大堆。5.3 Executor的GC参数ParallelGC为什么比G1更稳GC是影响Spark任务稳定性的大头Executor端JVM参数值得认真对待。大多数生产集群我会优先选ParallelGC而不是G1。原因在于Spark的任务特征是大对象多、吞吐优先、停顿敏感度低于在线服务ParallelGC在这种场景下Full GC频率更低吞吐量更稳。--conf spark.executor.extraJavaOptions-XX:UseParallelGC -XX:PrintGCDetails -XX:PrintGCDateStamps -Xloggc:/data/gc.log另外要留意Young区大小。Executor堆16G时我一般会把-Xmn放在4G左右也就是25%的堆比例。Young区小了大量短生命周期对象频繁进入Old区Old区很快就满触发Full GCYoung区过大又挤压Old区空间同样容易Full GC。需要结合GC日志观察。如果你非要用G1注意观察Mixed GC和G1 Full GC的节奏。G1在吞吐场景的优势不明显踩过几次坑之后我基本回归ParallelGC。当然这只代表个人生产经验如果你用的Spark版本和JDK版本不同建议用GC日志做A/B对照用数据说话。5.4 调优必须做A/B对照同一数据、同一时段、只改一个变量最后一条经验也是我觉得最重要的一条调优时要保证A/B对照的公平性。经常看到有人同时改动executor内存、分区数、并行度三四个变量然后说作业快了很多其实根本不知道是哪一项生效。更糟的是多租户集群里不同时间段的网络和调度资源波动很大同一个作业白天夜里跑可能差出两倍时间。我的做法是固定同一份测试数据、同一时间段、保证集群当前无其他队列抢占然后每次只改一个参数记录三个指标总耗时、GC耗时、Shuffle数据量。参数调整的方向依靠UI现象分析而不是拍脑袋。连续几轮下来结论才有说服力才能沉淀成可复用的调优模板。6. 一套可复制的Spark调优流程和参数速查把前面这些经验收敛一下得到的其实是一个很朴素的流程先看现象再定位瓶颈最后针对性调整。下面是我现在每接手一个新Spark作业都会走的流程你也可以直接拿来用。6.1 调优三步法现象、瓶颈、方案第一步用Spark UI确定瓶颈类型。是调度延迟高还是某个Stage的长尾Task还是GC时间占比大还是Shuffle数据量异常不同现象对应完全不同的优化方向这步错了后面全是无用功。第二步针对瓶颈类型选择方案倾斜就做加盐或广播Shuffle过大就优化Join策略和分区数GC严重就调JVM参数和序列化。第三步验证改动并回归记录对比数据。有一个真实的反面案例团队同学接手一个跑批作业上来就问executor给到64G行不行。我让他先跑一遍看UI结果发现瓶颈根本不在内存而在Shuffle阶段的数据居然有2TB原因是Join的某一侧缺少过滤条件。后来只是把过滤条件下推Shuffle数据量降到200GB任务从40分钟变到9分钟。这就是不定位瓶颈直接调资源的典型损失。6.2 常用参数速查表与调整方向基于这些经验我整理一份常用参数速查表方便你做线上调整时快速检索参数默认值调整方向适用场景spark.executor.memory1g按实际负载上调至8g-16g堆内内存不足、GC频繁或OOMspark.yarn.executor.memoryOverheadmax(384MB, 0.1*executorMemory)上调至1g-2g堆外OOM、JSON解析、Native内存任务spark.executor.cores12-3并行度不足但避免单Executor过载spark.sql.shuffle.partitions200按Shuffle量估算分区不均匀或分区过多导致小文件spark.default.parallelism无限与总核数匹配SparkContext级别的并行度控制spark.memory.fraction0.6建议保持0.6-0.7大量缓存时上调否则不下调spark.sql.autoBroadcastJoinThreshold10MB谨慎调整结合手动Hint小表Join大表场景spark.sql.adaptive.enabledtrue默认开启Spark 3.2以上spark.serializerJavaSerializer改为KryoSerializer大状态、复杂对象任务spark.shuffle.file.buffer32k上调至128k-512kShuffle落地文件过多、落盘压力大spark.reducer.maxSizeInFlight48MB上调至96MB网络带宽充足时减少Reduce端请求次数这些参数不是孤立生效的比如调整Shuffle分区数之后可能影响Task大小和GC表现建议每轮只动一个跑完一轮看ALL指标再动下一个。6.3 多租户环境下容易踩到的几个环境坑最后聊聊多租户集群里容易被忽视的坑。第一个是队列配额。同一个提交参数在A队列可能秒级获得资源在B队列要等资源释放表现差异非常大。我通常会在提交脚本里显式指定队列并且调优对比尽量固定在同一个队列里。第二个是YARN的节点本地性。如果某个作业经常从远端节点拉数据优先检查HDFS的块分布和YARN调度策略而不是Spark参数。第三个是并发作业之间的干扰尤其是动态资源分配时其他队列突然扩到海量Executor可能把你的作业挤到边角。做对比测试前先确认集群的全局负载是平稳的。6.4 几点实操体会如果要说这些年调Spark最深刻的心得大概是两句话一句是参数是下限代码是上限很多作业的根本问题在数据模型和代码写法参数只是帮你把代码跑得更顺另一句是调优不是一次性动作而是一套流程每个作业都应该沉淀自己的基线、现象库和参数模板。我个人在碰到任何新作业时会先把上面这套流程完整走一遍把三类指标记录下来放到一个共享表格里。这样下次同事报同一个作业变慢我可以直接对比历史数据判断是参数退化还是数据波动而不是重新从零开始分析。这套习惯救了我很多次建议你也试试。
返回列表