拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark作业性能调优:资源参数、代码优化与数据倾斜处理

Spark作业性能调优:资源参数、代码优化与数据倾斜处理

干了这么多年Spark,最常被问的问题就是:我的作业为什么这么慢?内存明明给够了,为什么一直在GC?数据量也就几千万,跑一个小时还出不来。

说实话Spark调优确实是个经验活,不是背几个参数就完事的。我前前后后调过几百个作业,从集群几百核到上千核的都碰过,踩过的坑比很多人跑过的任务都多。这篇算是我自己的一篇实战整理,把Spark调优从资源参数、代码写法到数据倾斜、监控排查串一遍,全是实操里验证过的东西。

很多人在网上找了一堆参数,什么spark.sql.shuffle.partitions设2000、spark.executor.memory直接往大了怼,然后发现该慢还是慢,甚至更慢。这不是参数不对,而是没搞清楚每个参数到底在干什么、跟你的集群和作业是不是匹配。我希望你看完这篇之后,能自己判断“我的作业该调哪里”,而不是照抄别人的配置。

1. 先搞清楚:Spark调优到底在调什么

1.1 一个任务跑得慢,问题藏在哪里

我之前带过不少新人,他们拿到一个慢作业,第一反应就是加内存、加核数,结果资源加了一倍,时间只是从40分钟变成35分钟。原因很简单:瓶颈根本不在资源上。

Spark作业的整个生命周期就那几段——读取数据、计算、Shuffle(中间数据落盘传输)、调度等待。你加内存,如果瓶颈是在某个数据倾斜的Task上,那加再多资源都白搭,因为倾斜的那个Task还是单点跑完才算完;如果瓶颈是Shuffle写磁盘太慢,你去调代码写法才有用。

所以调优的第一步不是动参数,是先看UI,定位瓶颈在哪。Spark自带的Web UI里,Stage、Job、Task的耗时、耗时占比、GC耗时都写得清清楚楚,先花十分钟把UI看懂,比调一百个参数都值。

我有一个习惯:拿到一个慢任务,先看两个地方。第一是Stage列表里每个Stage的耗时,第二是Timeline(时间线)上Task的分布情况。如果发现某个Stage里绝大多数Task几秒就结束了,就剩两三个Task跑了几分钟甚至报错,那不用犹豫,百分百数据倾斜。如果所有Task都慢,那就看GC时间,GC耗时长说明内存不够或者内存配比不合理。

1.2 调优的整体思路:四个维度一条线

我自己总结的Spark调优是四条线,你按顺序排查基本不会漏:

  • 资源维度:Executor数量、单Executor核数、内存大小、内存里堆内堆外比例。这是底层地基,地基本身要稳定,但不是越高越好。
  • 参数维度:并行度、Shuffle相关参数、序列化方式、动态分配开关。这是按作业特点微调的旋钮,要理解每个旋钮在转什么。
  • 代码维度:算子选择、缓存策略、广播变量、避免重复计算。这往往是收益最高、零成本的一部分,改代码比加机器靠谱。
  • 数据维度:倾斜处理、分区合理性、小文件治理。这部分最考验经验,也最容易拉开调优水平差距。

这四个维度不是割裂的,是一条因果链。你在代码里写了groupByKey,会导致Shuffle数据量巨大,Shuffle数据量大了,磁盘IO和网络IO就成了瓶颈,这个时候你再怎么调spark.executor.memory都没用,因为数据都在路上,根本没到内存里。反过来,你资源给得不够,分区数又顶得很高,那很多时间就消耗在线程切换和Task调度上了。

我自己排查问题时,习惯按“资源->参数->代码->数据”这个顺序来走一遍,每次只改一个变量,观察结果变化。如果一次改好几个参数,出了问题你根本不知道是谁的锅。

2. 内存与资源参数:最容易见效也最容易翻车的地方

2.1 Executor内存怎么分:堆内、堆外与Overhead

很多初学者以为spark.executor.memory就是Executor能用的全部内存,然后给Executor设个20G,跑起来OOM了就很懵。实际上这20G只是堆内内存,真正给Executor的总内存还要加上spark.executor.memoryOverhead,而且在Spark的内存管理模型里,堆内内存又被切成了三块。

先看那个经典的内存布局。堆内内存被spark.memory.fraction(默认0.6)分成两块,一块是执行和存储共享的内存池,专门干计算和缓存RDD这类活;另外40%是被“保留”的,主要放Spark内部对象、防止OOM的预留空间。在共享的那0.6里,又被spark.memory.storageFraction(默认0.5)分成执行内存和存储内存两块,中间可以互相借用。

另一个很多人容易忽略的是spark.executor.memoryOverhead,默认是executor.memory的10%,最小值384MB。这一块主要放JVM内部字符串、直接缓冲区、线程栈这些堆外开销。如果作业里大量用到collect拉数据到Driver,或者有比较大的Python/PySpark依赖,经常要手动往上加。

我再列一组我自己的经验值,这套配置在几十个生产集群上验证过,做数仓ETL、离线报表基本够用:

参数推荐值说明
spark.executor.memory8G-12G单Executor内存不建议超12G,超了GC会很痛苦
spark.executor.memoryOverhead2G-4G有大量序列化或Python作业时设高一点
spark.memory.fraction0.6-0.75内存充足可以调高,缓存/计算压力大
spark.memory.storageFraction0.5缓存RDD多的作业可以调到0.6以上

我踩过一个很典型的坑:有个作业缓存了一个几百MB的RDD,结果因为存储内存占比太低,缓存数据被Evict(驱逐)了,每次循环都重新计算一遍,任务直接慢了10倍。后来把storageFraction从0.5调到0.6,配合Kryo序列化,缓存的数据安安稳稳待在内存里,时间一下子从70分钟降到25分钟。

内存这里还有一个反直觉的点:不是堆越大越好。JVM堆超过32G就会失去指针压缩优化,可能反而更慢;堆太大导致GC停顿时间变长,Full GC一次可能好几秒。更合理的做法是增加Executor数量、保持单Executor内存适中,而不是堆一个超大内存的Executor。

2.2 并行度与Executor数量的配合

并行度这个概念,说白了就是让事情同时分给多少个人干。Spark里并行度的体现是分区数(Partition),分区数决定了Task的数量,Task是真正跑在Executor上的最小工作单元。

很多作业慢最直接的原因就是Task太少,比如数据量100GB,默认分区数可能只有几十个,那几十个Task把Executor占满之后,其他核都在空转。这时候你把并行度提上去,任务就会被切得更碎,分给更多核去跑。

但并行度也不是越高越好,这个真的经常有人搞反。我接过一个作业,跑数仓清洗,spark.sql.shuffle.partitions设了4000,结果每个Task处理的文件数据量很小,调度开销倒占了很大比例,整个作业光调度就花了七八分钟。后来改成根据数据量动态算分区数——大概是数据量的目标分区字节数除以128MB——作业整体从50分钟降到15分钟。

一个相对理想的公式大概是这样的:目标并行度 = 总CPU核数 * 2~3倍。为什么要多出来一点?因为Task不是绝对均匀的,有的Task快有的Task慢,多出来的并行度可以缓解木桶效应。如果是纯数据倾斜场景,那就另说,倾斜的Task再多也没用,后面专门讲。

spark.executor.cores这个参数也值得留意,它决定单个Executor能跑多少个Task。不建议设置太高,4-8个比较合适。太高了,同一个Executor上多个Task会争抢CPU和内存,尤其是GC和Shuffle的时候,互相拖后腿。

2.3 动态资源分配到底要不要开

默认情况下,Spark作业会申请你配置的所有Executor,比如你设了100个Executor,这些资源从作业开始到结束就一直占着,哪怕最后就剩一个Stage在跑最后一个Task,其余99个都在闲置。

动态资源分配就是解决这个问题的。开起来之后,集群会根据作业的实际负载自动调整Executor数量,空闲的还回去,繁忙的再加回来。对于业务峰谷明显、多人共用一个集群的场景,这个开关能省不少资源。

开启条件有三个:spark.dynamicAllocation.enabled=true、spark.shuffle.service.enabled=true,另外要是用了YARN做资源管理还需要External Shuffle Service配合。这背后的原理是,Executor退出之后,它持有的Shuffle中间数据需要有一个外部服务继续提供读取能力,否则其他Task在Fetch数据时就会找不到文件。

但我个人建议,如果是定时调度、数据量比较固定、资源本身不紧张的作业,其实可以不开动态分配,因为动态分配有抖动,Executor频繁启停本身就有开销。我遇到过动态分配把Executor缩到2个,结果后面一个Stage突然要处理的数据量很大,只能临时申请资源,等待资源的时间反而比固定配置更长。

所以我的判断标准很简单:作业负载比较稳定的,关掉动态分配,固定资源,确保每次运行的稳定性;作业负载随机性大的,比如即席查询、租户共享集群,打开动态分配。

3. 算子与代码层的调优:不同写法,性能差一个数量级

3.1 reduceByKey vs groupByKey:看起来一样,性能天差地别

先抛结论:能用reduceByKey的地方千万别用groupByKey。

这俩算子的区别得从Shuffle的原理说起。groupByKey是把所有key相同的value原封不动地Shuffle到同一个Task上,整个过程没有任何本地预聚合,数据量是原始的全量。举个例子,一个日志表里“用户A”这条key在文件里出现了一万次,那这就一万条记录全部通过网络传到下游。

而reduceByKey在Shuffle之前,先在每个Map端的Executor上做一次本地合并,比如同样是“用户A”在这个节点上出现了1000次,先在本地reduce成1条(或者几条),然后再Shuffle出去。这相当于数据量直接降了几个数量级。

我用一个实际项目验证过:一个网约车订单数据清洗作业,要按订单ID聚合金额和里程,数据量大概2亿条。groupByKey版本跑了一个小时,reduceByKey版本大概是11分钟。差距不是百分之几十,是几倍,就是因为Shuffle数据量完全不同。

类似的还有distinct和reduceByKey(_+_).map之类的选择,核心思路就是:想办法让数据在Shuffle之前瘦身。

3.2 广播变量与序列化优化:两个容易被忽略的隐形加速器

先讲广播变量。假设你有一张维度表,比如城市表,才几万行,要跟一个几十亿的事实表关联。如果你直接join,每个Task都会把这张表拉一份到自己的Executor上,网络开销很大,而且内存里每个Task都存一份副本,相当于白白浪费内存。

我在处理数据倾斜和大表小表Join的时候,第一个想到的就是broadcast join,把小的维度表广播出去。但这背后其实有一个阈值,spark.sql.autoBroadcastJoinThreshold默认是10MB,也就是小于10MB的表会自动走广播。但这里有个很坑的地方:这个10MB是“序列化后的大小”,不是表在DataFrame里看起来的大小。我遇到过一个配置表,看起来只有几万行,以为会自动广播,结果因为字段很多,序列化之后超过10MB,还是走了SortMergeJoin,跑了两个小时没出来。后来手动把阈值调大到50MB,作业6分钟搞定。

再说序列化。Spark默认的Java序列化很稳但很慢,还会把对象搞大。Kryo序列化虽然配置麻烦一点,但性能优势非常明显,尤其是Shuffle的时候,数据量少了,网络和磁盘开销都下来了。

启用方式很简单:spark.serializer=org.apache.spark.serializer.KryoSerializer,如果RDD用到了自定义类,还要注册一下。我通常还会配合spark.kryoserializer.buffer.max适当调大,默认64MB,任务里如果单个对象比较大,不够的话会报Kryo序列化失败。

3.3 RDD缓存与持久化的正确姿势

cache和persist是Spark里提效很常用的手段,但很多新手会无脑到处缓存,结果缓存的数据量太大,反而把执行内存挤没了。

先明确一个原则:只有一份RDD会被重复使用,或者会在多个Action里反复用到的时候,才值得缓存。如果它只被用一次,cache就是画蛇添足,既占内存,还因为计算父RDD增加了开销。

其次要选对存储级别。cache()默认是MEMORY_ONLY,数据以对象形式存内存;如果内存不够,部分分区缓存不了,之后要重新计算。我的经验是,如果数据能全部塞进内存,MEMORY_ONLY或MEMORY_ONLY_SER(序列化存储)最省CPU;如果可能出现内存不够,用MEMORY_AND_DISK_SER,宁可多写磁盘,也不要反复从头计算。

还有一个很多人没注意到的细节:缓存之后的RDD如果被unpersist掉,集群里资源才能真正释放。否则那些缓存的Block会一直赖在存储内存里,拖慢后面的作业。我在写调度任务的时候,都会在最后手动unpersist,释放掉不再用的缓存块。

4. 数据倾斜专项处理:90%调优场景的终极boss

4.1 怎么定位倾斜任务

数据倾斜是Spark慢作业的最常见元凶,也是面试里被问得最多的问题。四个字概括就是“木桶效应”:整个TaskSet的完成时间取决于最慢的那个Task。

怎么定位?还是老办法,打开Spark UI,看Stage里的Task耗时分布。正常情况Task耗时是接近均匀的,顶多有一些波动;倾斜的典型特征是“长尾”——绝大多数Task秒级完成,三五个Task跑了几倍甚至几十倍的时间。

再看数据倾斜的类型。一种是单个key特别多,比如“某个大商户”的订单占了全量数据的80%,加一个filter就能过滤掉的那种,倾斜很严重。另一种是key分布均匀但是值很大,比如某个key携带了一个巨大的数组字段,对网络和内存压力大。这两种的处理思路不太一样。

4.2 两阶段聚合与加盐写法

对于“少数key数据量特别大”的聚合场景,最经典的解法是“两阶段聚合”,江湖人称“加盐”或“打散”。

思路是这样的:因为某个key的数据量太大,导致所有数据都堆在一个Task上,那干脆把这个大key打散成多个key。具体做法是给key加一个随机前缀(通常是(key + "_" + rand.nextInt(10))),先让这些数据均匀分布到多个Task上做一次局部聚合,然后把前缀去掉,再做一次全局聚合。

两阶段聚合的效果我实测过,一个数据倾斜极其严重的订单统计作业,单个key占了全表40%的数据,没加盐之前那个Stage卡了40分钟出不来,加盐之后整体12分钟跑完。但要注意,加盐之后结果要能还原,比如你加了10个随机前缀,最后聚合时要stripSuffix把前缀去掉。

4.3 倾斜分区与大表关联的折中方案

如果倾斜发生在join场景,两阶段聚合有时候不管用,因为Join不像聚合那样可以简单分两步。这时候要分情况处理。

如果是大表Join小表,最佳方案就是前面的广播变量:把小表广播出去,避免Shuffle,直接从根上解决倾斜。我处理过一个用户行为表和用户维度表的Join,因为一个小部分用户的行为数据占了60%,没广播之前别的都跑完了,这一个用户的数据还在慢慢晃。改成广播Join后,整个作业时间压缩了80%。

如果是大表Join大表,还都倾斜,那就复杂一些。我的做法一般是这样:先把倾斜的key单独过滤出来,跟另一张表的对应key走广播Join,剩下的正常数据走正常Join,然后把结果union起来。这条思路虽然要写挺多代码,但在生产环境里很管用,也容易在各路调优方案中排上前几名的效果。

4.4 无法根治的倾斜如何兜底

有些倾斜是真的没法从业务上解决,比如数据本身分布就这么烂。这时候有一个应急兜底手段:把spark.sql.shuffle.partitions加大,比如从200调到1000或2000,让每个Task处理的数据量变小。虽然倾斜Task还是会慢,但因为它每个Task数据量变小了,整个Stage的瓶颈时间也跟着变小。

另一个兜底方案是开启spark.sql.adaptive.skewJoin.enabled(在Spark 3的AQE框架下),它会自动探测倾斜的分区,把倾斜的分区拆成更小的分区。说白了就是“让Spark自己帮你把大key切成多个小key”,省掉你手动加盐的步骤。

5. Shuffle与磁盘IO优化

5.1 Shuffle相关的关键参数

Shuffle是Spark作业里最耗资源、最容易出问题的一环,也是调优绕不开的重点。前面讲的加盐、广播、预聚合,本质上都是在“减少Shuffle的数据量”,但数据量已经固定的前提下,怎么让Shuffle过程本身更快,就要动一些参数了。

整理一下我常用的Shuffle参数:

参数默认值我的推荐值说明
spark.shuffle.file.buffer32KB64KB-128KBMap端写Shuffle数据时的缓冲区大小,调大减少磁盘IO次数
spark.reducer.maxSizeInFlight48MB96MBReduce端一次拉取数据的最大大小,调大减少拉取次数
spark.shuffle.io.maxRetries35-10Reduce端拉取失败后的重试次数,网络抖动时适度加大
spark.shuffle.io.retryWait5s10s-30s两次重试之间的等待间隔

这里特别说下spark.reducer.maxSizeInFlight。它的数值不是越大越好,太大了会导致每个Reducer一次拉很多数据到内存,反而增加GC压力。我之前为了减少拉取次数,一口气调到256MB,结果GC时间暴涨,刷新UI一看根本没有变快,后来调回96MB,时间和稳定性都对了。

spark.shuffle.io.maxRetries这个参数在网络不稳的集群上很重要。我遇到过某个混部集群,一到晚高峰YARN节点就时不时心跳超时,Shuffle Fetch失败率飙升,任务动不动就挂掉。把maxRetries调到10,retryWait调到30s,任务就很少因为这个挂了。

5.2 规避无谓的Shuffle

调参数是治标,代码层面少Shuffle才治本。

最常见的无谓Shuffle来自两大类:

  • 用了宽依赖算子但实际可以用窄依赖替代的场景。
  • 多次调用groupByKey聚合不同指标,但实际一次agg就能解决。

比如你有一个订单RDD,既要按城市算金额总量,又要按城市算订单数,如果你写了两次reduceByKey,就会产生两次Shuffle。改成一次Map生成结构化数据(城市,金额,数量),然后reduceByKey((a,b) => ...)一次搞定,Shuffle直接砍半。

还有一个频繁踩到的坑:repartition和coalesce乱用。repartition是宽依赖,会触发Shuffle;coalesce是窄依赖,多数情况下不触发。如果只是想把分区数调小,比如从几千降到几百,用coalesce就行;只有需要增加分区数或者让数据更均匀时再用repartition。

6. 集群部署与监控:调优之后可观测性才是关键

6.1 Spark History Server与EventLog

调优这件事必须建立在“能观测”的基础上。连一个作业哪里慢都看不到,调什么都是玄学。

最基础的一步:开启EventLog。配置spark.eventLog.enabled=true,再指定spark.eventLog.dir(比如HDFS上的一个目录)。这样作业运行结束后,事件日志会持久化下来,配合Spark History Server,随时可以回去查看历史作业的运行详情。

我之前接手过一个别人搭的集群,EventLog是关的,作业跑挂了就只能看YARN的日志,一堆乱码不说,Job级别的耗时、GC情况全都查不到。后来把EventLog和History Server打开,排查效率提升一个量级,很多原来“查不出来为什么慢”的问题,其实是每天固定某几个时刻在慢,看历史记录一眼就找到了。

部署History Server本身不复杂,把SPARK_HISTORY_OPTS配好,指定日志目录,启动start-history-server.sh就行。唯一要注意的是目录权限,如果目录没有写权限,EventLog写不进去,作业会直接报错,这个我踩过坑。

6.2 常用监控指标和工具选择

Spark UI本身其实是最好用的监控工具,没有之一。里面有几个关键指标是我每次必看的:

  • Completed Stages / Active Stages:这个作业整体在哪些阶段耗时最多。
  • Task Deserialization Time / Serialization Time:如果占比很高,说明序列化方式可能没调好。
  • Shuffle Read / Write Size:看数据在Shuffle阶段到底传输了多少。
  • Executor GC Time:GC时间占比超过10%就要重视了,否则说明内存配置有问题。
  • Task的Locality Level:PROCESS_LOCAL是最理想的,如果是NODE_LOCAL甚至RACK_LOCAL,说明数据本地性不好,网络传输成本偏高。

除了Spark UI,我们在生产环境还配合用了Ganglia和Prometheus。Ganglia看集群整体资源水位,判断是不是有人抢占资源;Prometheus+Grafana用来画作业运行趋势图,做告警。比如Shuffle Fetch失败率突增、Executor OOM次数突增这类,配置好告警规则能提前发现问题。

如果集群规模大到用YARN管理,还有一个yarn application -status的快速查看方式,能快速定位Application提交状态和资源使用情况。不过说实话,日常调优用Spark UI基本够用,那些外部工具是在集群有持续性性能问题的时候才需要上的。

7. 常见问题与排查技巧实录

7.1 典型报错与解决路径

先把自己常碰到的几类问题整理成了速查表,按“现象->原因->解决”的顺序来写,你可以直接拿去对照:

现象大概率原因首选排查方案
Executor Lost / Container killedExecutor内存不足,或Overhead不够调大executor.memory或memoryOverhead,看YARN日志确认
java.lang.OutOfMemoryError: Java heap space堆内内存不足调大executor.memory,检查是否有倾斜,检查缓存是否滥用
Shuffle fetch failed网络不稳,或Executor被回收调大maxRetries和retryWait,检查是否有Executor OOM
Job aborted due to stage failure:Serialized task太大广播变量或Task携带数据太大调大spark.rpc.message.maxSize,或优化广播变量
GC耗时占比超高内存比例不合理,对象过大调整memory.fraction,换Kryo序列化,检查缓存数据量
所有Task都慢且耗时平均集群资源不足,或并行度太低增加Executor,调大default.parallelism/partitions

这里要说一个很多人忽视的点:Serialized task太大这个报错,很多时候不是Task本身大,而是在collect时Driver端把大量数据序列化收回来,超过了RPC消息大小限制。解决方式是别用collect拉大数据,改写take或分批拉取,或者调spark.rpc.message.maxSize到1024MB兜底。

7.2 快速定位慢作业的思路与心得

我一直跟团队里的人讲,排查慢作业别上来就动手调,要先画出来、看明白。我自己的固定套路是:先看作业整体耗时最长的Stage,再定位到具体的Task类型(读HDFS?Shuffle Read?还是CPU计算?),最后看数据分布和GC占比。

举个例子,之前有个作业每天凌晨3点跑,数据大概几十GB,原本一切正常。突然有一天跑到一半卡住,最后超时失败。我打开Spark UI,发现某个Stage里有一批Task在疯狂Retry,Shuffle Read的数据量比前一天翻了好几倍。最后查下来发现,数据源方前一天改了分桶策略,输出了一堆超大的文件,数据分布烂了,而且刚好撞到之前设的Shuffle分区数不够,每个Task拉的数据量暴涨。

这个问题的解决其实很简单:把spark.sql.shuffle.partitions从200调到400,问题当天就解决了。但如果没有UI定位,先改内存、再改GC,可能折腾一晚上都找不到原因。

所以调优的经验再多,最后都要落到“会看、会定位”这四个字上。

8. 一点点经验总结

我不太喜欢给人列“标准答案”,因为每个集群、每个作业的脾气都不一样。但我在实际项目中确实验证过一个通用结论:调优的收益排序,数据倾斜处理大于代码优化大于内存参数大于并行度参数。代码层面的改写往往能带来几倍甚至十几倍的提升,参数调优的收益通常也就百分之几十。

如果你刚接触Spark调优,我建议先别急着改配置,把一个作业的UI彻底看懂:哪个Stage慢、哪个Task倾斜、GC时间占多少、Shuffle传输了多少,全搞明白。这个过程看着枯燥,但真的是调优最值钱的基本功。

最后再分享一个小技巧:每次调优只改一个参数或者一段代码,记录前后的运行时间和资源使用。很多人在几个优化点之间来回摇摆,最后根本不知道哪个改动起了作用。用这种“单变量试验”的思路,你会比很多人都走得快。

返回列表