写Spark调优这个话题我最怕被问“你给个参数模板我直接抄”。说实话调优这件事抄是抄不出来的不搞明白底层的资源模型和执行机制给你一套看起来合理的配置放到你的集群上反而可能更慢。但把底层的逻辑讲透把常见的坑提前排掉却是可以写成文字、反复拿出来对照的。这篇东西就是我这些年在大数据一线折腾Spark被数据倾斜折磨过、被Executor OOM教育过、也被G1参数救回来之后沉淀下来的一套实战经验。不管你是在做集群搭建、跑网约车数据分析还是刚接触Spark开发只要你用的是Spark处理上亿条数据这篇文章都值得你从头看一遍。我会把资源参数、代码习惯、数据倾斜、内存GC这几个最容易出问题的地方挨个拆开每一条都是我在生产环境里验证过的有的带命令有的带代码尽量让你看完就能拿去用。1. 先搞清楚Spark调优到底在调什么1.1 调优的本质资源、并行度、内存和GC之间的博弈很多人一上来就调spark.executor.memory和spark.executor.cores这没错但容易陷入一个误区以为内存调大、核数调多作业就一定会变快。Spark的调度模型是把一个DAG切分成多个Stage每个Stage内部根据RDD的分区数决定并发Task的数量。Task的运行速度和两个东西强相关一是计算本身的复杂度二是Task所在Executor能拿到的CPU和内存资源。资源给得再多如果并行度不够跑起来就是“大马拉小车”资源给得紧并行度又拉得很高GC和调度开销反而会把性能吃掉。所以调优的第一步不是调参数是理解你作业的“瓶颈形态”是数据倾斜导致的慢是shuffle文件太多导致的网络IO还是单Task内存开销太大导致的频繁Full GC。1.2 从一次数据倾斜事故说起问题表象 vs 根因我之前接过一个网约车日志清洗作业每天的数据量大概三个亿的订单点击日志。表面现象是某个Stage运行了40分钟其他Stage都是两三分钟完事。看Spark UI发现某个Task处理了600万条数据其他Task只有几千条明显是key分布极不均匀。这个案例表象是“某个Task特别慢”但根因是业务主键里包含大量空值和个别热点司机ID。如果你不定位到具体key分布只去加Executor内存那就是浪费资源热点Task一样跑不动。所以我的习惯是遇到慢任务先打开Spark UI的Stage详情页看Summary Metrics里面的Max时间、GC时间、Shuffle Read/Write大小再直接截取数据用groupBy看top key。定位不准后面所有调优都是盲打。2. 资源参数配置第一层也是最容易被无脑抄的参数2.1 一套建议的初始参数模板生产环境常见配置参数这东西不能脱离集群说。假设你有一台集群每台物理机128GB内存、32核YARN为每个NodeManager分配100GB、26核下面这组参数是我在类似规模下比较稳的起步模板适合批处理场景Streaming要按延迟需求再压参数名设置值一句话说明spark.executor.cores4每个Executor占4个核避免单个Executor独占整机方便并发运行多个作业spark.executor.memory12g堆内存预留 yarn overhead 后估算spark.executor.memoryOverhead2g堆外内存跑Python和复杂C算子时尤其重要spark.driver.memory4gDriver主要存任务元数据和collect结果别给太小spark.default.parallelism300起步值实际按CPU核心数×2~3来定spark.sql.shuffle.partitions300Spark SQL shuffle后默认分区数同上spark.memory.fraction0.8留给统一内存的堆比例剩下给用户代码spark.memory.storageFraction0.5执行和存储各占统一内存的一半spark.yarn.executor.memoryOverhead2gYARN模式下堆外内存配置和上一个保持一致即可注意以上模板不是让你无脑抄。每个作业的输入大小、shuffle规模、udf复杂度都不同真正上线前必须在测试环境跑一遍小数据量加压测。2.2 关键参数背后的计算逻辑很多人不知道spark.executor.memory和spark.executor.memoryOverhead加起来才是YARN分配给容器的总量。比如你设置executor.memory12goverhead2g那单个container的物理内存就是14g如果你忘了overheadYARN会默认按堆内存的10%给结果往往不够用容器直接被NodeManager杀掉报Container killed by YARN for exceeding memory limits。这算是新手最常见的一个坑。再说并行度我更倾向于用spark.default.parallelismtotal_executor_cores × 2 ~ 3来估算。假设你有20个Executor每个4核总核数80并行度给160到240比较合理。并行度太高每个Task处理的数据量小调度和序列化开销占比升高太低CPU又跑不满。另外spark.sql.shuffle.partitions在Spark 3.2之后默认是200但实际生产里shuffle分区不一定越多越好。我曾经试过一个countByKey作业把分区数从200调到800反而慢了20%因为每个分区写入磁盘的随机IO更多了。2.3 我的参数调整心法先看监控再动刀调参数最忌讳一次改五六个然后重新跑一遍看结果因为出了问题你根本不知道该回退哪个。我在实际项目里的做法是先跑一次baseline把Spark UI里的“Completed Stages”页面所有指标截下来重点看Shuffle Read Size和Shuffle Spill Memory。如果发现spill很大说明Executor内存不够或分区数据不均如果发现某个stage的Task分布极不均衡先处理倾斜再看要不要调内存。参数调整一次只动一个变量比如先调并行度记录跑批时间再动内存再记录。整个过程最好配一个简单的脚本把每次改动和运行时长记成一张表时间长了你就有了自己的“调优基准库”。3. 代码层面的调优技巧不换参数也能提速50%的写法3.1 RDD还是DataFrame选错API后期全是泪有很多从RDD迁移过来的老代码习惯用mapPartitions去做复杂解析。不是说RDD不能写但Spark SQL的Catalyst优化器和Tungsten执行引擎在DataFrame上的优势太明显了谓词下推会让过滤条件尽早执行列式存储会大幅减少IO向量化处理让单条数据处理的CPU指令更高效。我见过最简单的例子同样从读取Parquet文件到filter再到groupByDataFrame API比RDD API快了近60%代码量还少了一半。下面这段对比很直观// RDD方式序列化开销大无法走Catalyst优化 val result df.rdd .filter(row row.getAs[String](city) 上海) .map(row (row.getAs[String](driver_id), 1)) .reduceByKey(_ _) // DataFrame方式自动优化可读性和性能都更好 val result df .filter(col(city) 上海) .groupBy(driver_id) .count()你如果在生产环境维护过一堆RDD代码一定会明白“用DataFrame不是炫技是给自己省麻烦”这句话。哪怕遇到很复杂的业务逻辑也可以先用DataFrame处理主体再落到小数据集上用RDD跑精细逻辑。3.2 慎用groupByKey和distinctshuffle是万恶之源Spark里shuffle意味着数据要跨节点传输和落盘。groupByKey把每个key对应的所有value在map端不合并完整地传到下游数据量巨大。而reduceByKey在map端先做一次本地聚合再传结果网络IO能小一个数量级。举个非常实在的例子统计每个城市的打车订单数// 不推荐所有原始value都会跨网络传输 rdd.groupByKey().mapValues(_.size) // 推荐map端先合并宝贵 rdd.reduceByKey(_ _)再来说distinct。distinct底层就是map reduceByKey耗时不短。如果只是去重统计数量最好用approx_count_distinct底层是HyperLogLog这类近似算法误差在1%以内性能提升明显。比如统计UV用DataFrame.groupBy(dt).agg(countDistinct(user_id))在十几亿数据上可能要跑十几分钟换成approx_count_distinct(user_id, 0.05)往往两三分钟就完事而且这种精度在业务报表里完全够用。3.3 广播变量、缓存级别和序列化的选择有大表join小表时务必用Broadcast Hash Join。Spark 2.0之后会自动针对10MB以下默认spark.sql.autoBroadcastJoinThreshold的表做广播但很多时候小表实际有几百MB阈值要先调一下或者手动broadcast提示import org.apache.spark.sql.functions.broadcast val result largeDf.join(broadcast(smallDf), Seq(city_id), left)广播变量和缓存是两个概念。缓存是把中间结果放到内存或磁盘避免重复计算。要注意的是cache()是懒执行的必须在count()或写入操作后才真正触发缓存。很多新手写了df.cache()然后马上打印日志发现下次用还是慢就是因为缓存还没物化。序列化框架上Scala作业建议直接开启Kryospark.kryo.registrationRequiredfalse spark.kryoserializer.buffer.max512m spark.serializerorg.apache.spark.serializer.KryoSerializerJava序列化虽然兼容性好但数据体积大、序列化慢。Kryo在shuffle量大的作业里能明显降低GC和网络压力。3.4 动态分区和文件小文件问题的处理生产环境最常遇到的一个问题是写入Hive表时生成了一堆小文件比如几万个几十KB的文件下游查询拉取时NameNode和磁盘IO直接爆炸。解决小文件有几个手段一是设置spark.sql.shuffle.partitions不要过大让每个任务输出文件数可控二是写之前用coalesce()合并分区三是开启Spark的动态分区写入优化spark.sql.optimizer.dynamicPartitionPruning.enabledtrue spark.sql.adaptive.enabledtrue spark.sql.adaptive.coalescePartitions.enabledtrue动态分区裁剪DPP的作用是当事实表join维度表且以维度表字段为分区键时Spark能提前跳过无关分区大大减少扫描数据量。这个功能在Spark 3.0以后成熟如果你的作业还在Spark 2.4之前建议尽早升级收益远比调那些参数香。4. 数据倾斜实战表现、定位与四种解决姿势4.1 倾斜的症状与诊断命令倾斜的典型表现是同一Stage里大部分Task几秒跑完一两个Task要跑几十分钟而且Executor的GC时间飙升。定位方法很简单Spark UI的Stage详情页点击某个Executor看到Shuffle Read Records特别大或者看某个Task的“Input Size / Records”是其他任务的几百倍基本就锁定了倾斜的Task。接着要定位到具体Key可以把疑似数据拉出来跑个countdf.groupBy(col(key)).count() .orderBy(col(count).desc) .limit(20) .show(false)如果发现某个key的count异常大比如占比超过10%就是它在捣鬼。还有一种隐蔽倾斜是join的关联字段里有大量默认值比如“unknown”或“0”这些空值会被HashPartitioner集中到同一个分区。4.2 加盐/两阶段聚合解决group by倾斜加盐Salting是我用得最多的倾斜解法思路就是先把热点key打散加上一个随机后缀做第一轮局部聚合再去掉后缀做第二轮聚合。// 假设发现 hotKey 是热点 val salted df .withColumn(salt, when(col(key) hotKey, (rand() * 50).cast(int)).otherwise(0)) .withColumn(salted_key, concat(col(key), lit(_), col(salt))) val firstAgg salted .groupBy(salted_key) .agg(sum(value).as(partial_sum)) val finalResult firstAgg .withColumn(original_key, regexp_replace(col(salted_key), _\d$, )) .groupBy(original_key) .agg(sum(partial_sum).as(total_sum))这里加盐的粒度要控制好加盐数太多第一轮聚合会多一次shuffle太少热点key还是大。我一般对热点key加30到50个随机数然后观察第二轮聚合的Task分布是否均衡。两阶段聚合只适合groupBy agg的场景如果是join的倾斜不能用这个办法。4.3 大表join小表用广播大表join大表用拆分join倾斜是另一种常见病。两个大表join时如果关联key分布不均那么key比较多的分区会积压。一个方案是把大表拆成一个“热点部分”和一个“正常部分”。热点部分单独加随机前缀另一张表复制多份同样加前缀这样把热点key分散到多个Task正常部分保持原始key join。最后用union把两个结果合并。val hotKeysDf dimDf.filter(col(key).isin(hotKeyList: _*)).withColumn(prefix, explode(array((0 to 49).map(lit(_)): _*))) val hotFact factDf .join(broadcast(hotKeysDf), Seq(key), inner)这个方法写起来略繁琐但确确实实能解决80%的join倾斜。另外如果一张表有大量key为null最好先过滤掉null或者单独用非null的key去join否则null全堆到一个任务里也会被“假倾斜”骗进去改了半天发现只是空值太多。4.4 调整并行度和重新分区coalesce/repartition有时候倾斜并没有那么严重只是默认并行度不够把spark.sql.shuffle.partitions调大就能把数据分散到更多Task。注意调整后要观察单个Task的处理时间如果Task数多了但每个Task处理的数据量变小耗时自然下降。但也要防止从一个极端走到另一个极端分区数太大shuffle文件变成成千上万个零碎小文件下游reduce时反而因为IO拖慢。重分区操作上coalesce会尽量避免shuffle只合并分区repartition则会全量shuffle。我见过有人在写完文件前用repartition(1)导出单文件数据量一大就是灾难。更合适的做法是先coalesce(5)写出5个文件或者用partitionBy按业务键分目录写。5. 内存管理与GC调优最容易出问题也最容易忽略5.1 Spark内存模型再回顾Spark内存主要分为Execution和Storage两大区域加上预留和用户代码的空间。Execution内存用于shuffle、join、aggregate的缓冲区Storage内存用于缓存RDD、DataFrame。两者是动态占用的也就是说当Execution不够用时可以侵占Storage空闲部分反之亦然。这个设计的初衷是提高内存利用率但也带来一个坑如果你设置了df.cache()又没有给Storage留够空间缓存可能被频繁淘汰导致重复计算。虽然Spark UI里缓存显示的是“in memory”但你Task重跑一遍就会看到“block not found”这说明缓存实际被清掉了。我的建议是如果缓存的数据会被复用很多次就加大spark.memory.storageFraction到0.6甚至0.7如果作业是纯计算、shuffle很重反而可以调低到0.3让更多内存给Execution。5.2 Executor OOM与GC频繁怎么定位OOM不一定真的是数据量大更多时候是“计算过程中某个Task把内存爆了”。比如rdd.mapPartitions里有一个超大对象比如把整个文件load进内存再逐行处理那一个Task就能塞满整个Executor。定位方式除了看堆栈和Spark UI还可以用MySQL和tracer这类的线上工具但我个人经验是“自查代码”永远比查监控更快凡是遍历大集合时用了.toList或算出中间结果后没有及时清理都是OOM的温床。GC频繁的表现是Active Task数不高但Executor页面显示GC时间很长比如一个Task运行10分钟其中8分钟都在GC。这种情况要先看是不是堆内存设置得太小再看是不是分配给Execution的缓冲对象生命周期太长。如果数据本身正常调整并行度让每个Task的中间结果变小往往就能缓解。5.3 GC调优实用参数G1等JDK 8以上的Spark任务可以切换到G1垃圾收集器配合以下几个字段对大堆Executors效果要比默认的CMS好不少--conf spark.executor.extraJavaOptions-XX:UseG1GC -XX:PrintGCDetails -XX:PrintGCDateStamps --conf spark.executor.extraJavaOptions-XX:InitiatingHeapOccupancyPercent45 --conf spark.executor.extraJavaOptions-XX:G1HeapRegionSize16mInitiatingHeapOccupancyPercent代表G1触发并发标记周期的堆占用阈值默认45%不要调太高否则混合GC会延迟容易产生长时间停顿。G1的Region变大后大对象分配会减少触发Humongous Allocation的几率。以下表格是我在不同场景下用过的组合场景堆大小GC参数组合效果普通SQL批处理shuffle中等12gG1IHOP45Region16m稳定Full GC很少缓存大量维度表20gG1IHOP55Region32m减少STW缓存不再频繁换出机器学习迭代计算24gParallelGC默认不使用G1迭代计算中G1的并发标记反超不如Parallel稳定关于GC我的体会是不要迷信G1。很多迭代型算法比如需要频繁扫描中间数据的ALS类作业使用ParallelGC反而比G1更稳关键还是要在测试环境跑一版对比。6. 我总结的调优排查流程与日常巡检清单6.1 排查SOP五步法我把一次Spark性能问题从发现到解决的流程固定成五步新老手都能少走弯路看Spark UI点开Application页面找到耗时最长的Stage确认瓶颈是某些Task还是全体Task都慢。看输入数据读一下数据量、分区数、是否有明显热点Key这一步可以直接写一条简单的groupBy count命令跑出来。看Shuffle情况重点看Shuffle Spill如果把数据溢写到磁盘说明Executor内存不足或并行度不合理。微调资源参数一次只改一个变量记录变化不盲调。分析执行计划用df.explain(true)看是不是走了全表扫描或者join顺序能不能被优化器调整。这个SOP帮我在不少项目里快速定位过问题比如最麻烦的数据倾斜往往在第2步就发现了根本不用改参数。6.2 日常监控指标看哪些除了Spark UI和YARN ResourceManager我平时固定看这几项指标Shuffle Spill (MemoryDisk)溢写是内存和并行度失衡的信号。GC Time单个Task的GC时间超过总运行时间的10%就要警惕。Input / Shuffle Read Size这个决定IO压力能侧面反映分区是否合理。Executor的Active Task数如果长期等于Executor的核数说明有等待可能并行度不够或资源被其他作业抢占。HDFS文件数如果你有写入Hive的分区表频繁检查输出文件的数量和大小。6.3 几个容易误判的“优化”陷阱第一盲目增大广播阈值。把autoBroadcastJoinThreshold调到几百MB可以解决一部分小表join慢但Driver要把广播数据分发到每个Executor如果Executor数量多广播本身就会很慢甚至OOM。第二为了“一看就很专业”而加了一堆参数比如调一堆spark.sql.adaptive参数但没升级Spark版本实际上Spark 2.x根本不支持等于白调。第三你以为df.cache()加快了速度但缓存序列化耗时可能比重新计算还高尤其是数据被全局扫描一遍就扔掉的场景。提示调优的目的不是让所有参数都“有值”或者“最大”而是在你的数据规模、集群配置、业务需求三者之间找到平衡。如果某个参数改了没有明显变化就把它恢复原样保持集群配置的整洁。最后说一点我个人的体会。做Spark调优这几年最能提升我判断力的不是在社交媒体上看别人分享的“性能提升三倍”案例而是自己把Spark UI打开一个Stage一个Stage地看把每个Task的耗时、GC、Shuffle记录都核对清楚。调优是慢功夫很多时候你改一个参数要跑几十分钟才知道结果但只要你肯坚持记log、记参数、记现象慢慢就会形成一套自己的调优方法论。我的建议是第一次接触某个Spark作业时先不忙优化跑通、看UI、记录基线然后再动刀。稳比快重要一万倍。