
文章目录一、课前导读二、学习目标三、核心理论知识点四、原理通俗讲解4.1 资源调优让你的钱花在刀刃上4.2 内存管理Spark的“厨房”怎么安排4.3 数据倾斜少数“钉子户”拖垮整个任务五、重点概念拆解5.1 核心资源配置参数详解5.2 内存管理参数5.3 序列化优化5.4 Shuffle优化参数5.5 SQL优化要点5.6 数据倾斜治理方案六、易错点避坑6.1 盲目增加Executor数量6.2 忽略数据本地性6.3 动态资源分配关闭导致资源浪费6.4 过早使用collect()或toPandas()6.5 忽视小文件问题七、完整实战案例7.1 初始代码未优化7.2 优化后代码7.3 SQL调优示例7.4 监控与诊断八、代码逐行解析8.1 资源配置参数8.2 AQE关键参数8.3 加盐打散实现九、业务场景落地应用9.1 场景一大表Join大表优化9.2 场景二多维度聚合中的倾斜9.3 场景三流计算中的数据倾斜9.4 场景四机器学习特征工程中的倾斜十、常见报错排查10.1 Container killed by YARN for exceeding memory limits10.2 FetchFailedException10.3 OutOfMemoryError: GC overhead limit exceeded10.4 org.apache.spark.shuffle.MetadataFetchFailedException十一、本节课知识点总结调优参数速查表数据倾斜治理决策树性能优化流程十二、课后思考作业作业一理论理解题作业二代码实践题作业三场景应用题作业四拓展研究《20节课 PySpark 从入门到精通》系列课程导航一、课前导读经过前面18节课的学习你已经掌握了PySpark的绝大多数核心技能从RDD、DataFrame到Spark SQL从离线数仓到实时流计算。你写的程序已经能够正确地运行并输出期望的结果。但是在真正的生产环境中“能跑”只是第一步“跑得快”才是关键。你可能遇到过这样的问题明明集群有几百个核心但任务执行时却只有一个Task在运行其他人都在等待。一个小小的Join操作居然触发了TB级的Shuffle磁盘被撑爆。99%的Task都在几秒内完成却有1%的Task卡了半个小时拖垮了整个作业。开启动态资源分配后Executor不断申请和释放浪费了大量资源。这些问题归根结底都是性能问题。Spark的性能调优是一个系统工程涉及资源配置、数据序列化、内存管理、并行度设置、数据倾斜处理、SQL优化等方方面面。很多开发者只知其然不知其所以然盲目调整参数往往事倍功半。本节课将系统性地讲解PySpark性能调优的完整知识体系。我们会从资源参数调优Executor、内存、并行度开始然后深入Spark SQL优化Catalyst、AQE、分区剪枝最后重点攻克大数据领域最棘手的问题——数据倾斜提供多种行之有效的治理方案加盐打散、广播Join、自定义分区器等。学完这节课你将具备诊断性能瓶颈、主动优化Spark作业的能力成为一名真正的PySpark专家。二、学习目标完成本节课的学习后你将能够掌握核心资源配置参数合理设置Executor数量、内存、核心数、并行度等理解内存管理模型区分堆内内存、堆外内存调优spark.memory.fraction等参数优化序列化使用Kryo序列化替代Java序列化提升网络传输效率应用Spark SQL优化技巧利用AQE、动态分区剪裁、Join策略选择等根治数据倾斜识别倾斜场景使用加盐、广播、自定义分区器等多种手段解决优化Shuffle减少Shuffle数据量使用map端预聚合、调整shuffle分区数使用监控工具通过Spark UI、事件日志、Metrics系统定位性能瓶颈综合调优案例完成一个从慢到快的真实业务优化过程三、核心理论知识点调优维度关键点常用参数资源参数Executor数量、内存、核心--num-executors,--executor-memory,--executor-cores并行度分区数、任务并发度spark.default.parallelism,spark.sql.shuffle.partitions内存管理统一内存模型、存储内存、执行内存spark.memory.fraction,spark.memory.storageFraction序列化Kryo vs Javaspark.serializer,spark.kryo.registratorShuffle优化数据压缩、文件合并、缓冲区大小spark.shuffle.compress,spark.shuffle.file.bufferSQL优化AQE、动态分区剪枝、Join策略spark.sql.adaptive.enabled,spark.sql.autoBroadcastJoinThreshold数据倾斜加盐、广播Join、自定义分区器业务逻辑处理监控Spark UI、History Server、Metricsspark.eventLog.enabled四、原理通俗讲解4.1 资源调优让你的钱花在刀刃上集群资源是有限的就像你有一笔预算需要合理地分配给各个项目。如果给一个任务分配了太多Executor其他任务可能无资源可用如果分配太少任务又跑得慢。关键是要平衡并行度和资源利用率。Executor vs Core vs Memory每个Executor是一个JVM进程负责执行Task。过多的JVM进程会增加调度开销。每个Core可以并发执行一个Task。一个Executor的Core数决定了它可以同时跑多少个Task。内存大小决定了能缓存多少数据、能容纳多大的Shuffle数据。调优黄金法则Executor总数 总Core数 / 每个Executor的Core数每个Executor内存建议4-8GBCore数建议3-5个。4.2 内存管理Spark的“厨房”怎么安排Spark Executor的内存就像一个大厨房分为几个区域执行内存炒菜区用于Shuffle、Join、排序等操作需要频繁“颠勺”。存储内存冰箱用于缓存RDD、DataFrame、广播变量。用户内存存放用户数据结构和UDF内部对象。Spark 1.6的统一内存管理模型允许执行内存和存储内存互相借用如果对方空闲提高了内存利用率。但如果配置不当可能导致频繁的垃圾回收GC或溢写磁盘。4.3 数据倾斜少数“钉子户”拖垮整个任务数据倾斜是最常见也最头疼的性能问题。想象一下你要统计每个省份的人口但某个省份比如广东的人口远多于其他省份那么负责处理广东数据的那个Task就会非常慢其他Task早已完成只能干等。Spark中的倾斜通常发生在Shuffle阶段相同key的数据汇聚到同一个分区导致该分区的数据量巨大。解决思路就是打散给倾斜的key添加随机前缀让它们分散到多个分区然后分两步聚合。五、重点概念拆解5.1 核心资源配置参数详解spark-submit\--masteryarn\--deploy-mode cluster\--num-executors20\# Executor数量YARN模式--executor-cores4\# 每个Executor的CPU核心数--executor-memory 8g\# 每个Executor的堆内存--confspark.executor.memoryOverhead2g\# 堆外内存占总内存10-20%--confspark.driver.memory4g\# Driver内存--confspark.default.parallelism100\# 默认并行度Shuffle后分区数--confspark.sql.shuffle.partitions200# SQL中Shuffle分区数调优建议Executor内存不要超过32GB否则JVM GC暂停时间过长。推荐8-16GB。Executor的Core数通常取3-5避免超过5导致HDFS I/O吞吐瓶颈。并行度分区数应设置为总Core数的2-3倍让资源充分利用且减少调度开销。堆外内存至少分配10%尤其在使用UDF或处理大字段时。5.2 内存管理参数参数默认值含义调优建议spark.memory.fraction0.6统一内存占堆内存的比例剩余0.4为用户内存若缓存多可调高到0.7-0.8若UDF复杂保持0.6spark.memory.storageFraction0.5存储内存占统一内存的比例缓存多则调高到0.6Shuffle多则调低到0.4spark.memory.offHeap.enabledfalse是否使用堆外内存对内存敏感且大量缓存时可开启需配置offHeap.sizespark.sql.adaptive.enabledtrue(3.x)自适应查询执行强烈建议开启5.3 序列化优化默认使用Java序列化org.apache.spark.serializer.JavaSerializer性能较差。建议使用Kryo序列化sparkSparkSession.builder \.config(spark.serializer,org.apache.spark.serializer.KryoSerializer)\.config(spark.kryo.registrationRequired,false)\.config(spark.kryo.unsafe,true)\.getOrCreate()Kryo可以将数据序列化为更紧凑的字节数组速度也更快。如果知道要序列化的类可以注册以进一步提升性能。5.4 Shuffle优化参数参数默认值说明调优spark.shuffle.compresstrue是否压缩shuffle输出开启使用snappyspark.shuffle.file.buffer32k写文件缓冲区大小可增大到64k-128kspark.reducer.maxSizeInFlight48m同时拉取的数据量网络好可增大到96mspark.shuffle.sort.bypassMergeThreshold200小于此分区数时使用 bypass可增大减少排序spark.shuffle.partitions200SQL shuffle分区数根据数据量调整每分区100-200MB5.5 SQL优化要点开启AQEAdaptive Query ExecutionSpark 3.x默认开启能够动态合并shuffle分区、动态调整Join策略、处理数据倾斜。spark.conf.set(spark.sql.adaptive.enabled,true)spark.conf.set(spark.sql.adaptive.coalescePartitions.enabled,true)spark.conf.set(spark.sql.adaptive.skewJoin.enabled,true)spark.conf.set(spark.sql.autoBroadcastJoinThreshold,10485760)# 10MBJoin策略选择小表10MB应使用广播JoinBroadcastHashJoin避免Shuffle。中等表10MB-100MB可考虑广播。大表对大表使用SortMergeJoin。分区剪枝在查询中务必加上分区字段的过滤如WHERE dt2024-01-01。列剪枝只SELECT需要的列避免读取无关数据。5.6 数据倾斜治理方案方案一加盐打散两阶段聚合适用于groupBy/reduceByKey等聚合操作。# 给key加随机前缀0-9df_with_saltdf.withColumn(salted_key,concat(col(key),lit(_),(rand()*10).cast(int)))# 第一次聚合局部partialdf_with_salt.groupBy(salted_key).agg(sum(value))# 去掉盐第二次聚合resultpartial.withColumn(original_key,split(col(salted_key),_)[0])\.groupBy(original_key).agg(sum(sum(value)))方案二广播Join适用于大小表Join倾斜的key存在于大表但小表可以广播。frompyspark.sql.functionsimportbroadcast resultlarge_df.join(broadcast(small_df),key)方案三拆分倾斜key将倾斜的key单独处理与普通key分开Join后再合并。# 识别热点key例如出现次数1000hot_keysdf.groupBy(key).count().filter(count 1000).select(key).collect()# 分离热点和非热点数据hot_dfdf.filter(col(key).isin([k[0]forkinhot_keys]))normal_dfdf.filter(notcol(key).isin([k[0]forkinhot_keys]))# 对热点数据用广播Join或其他方式hot_resulthot_df.join(broadcast(dim_df),key)normal_resultnormal_df.join(dim_df,key)# 合并resulthot_result.union(normal_result)方案四自定义分区器对于RDD或DataFrame的repartitionByRange可以自定义分区逻辑将热点key分散。六、易错点避坑6.1 盲目增加Executor数量增加Executor会导致调度开销增大且每个Executor都会申请内存可能超出集群容量。应先增加Core数或并行度而非Executor数。6.2 忽略数据本地性如果Task被调度到没有数据的节点需要从远程拉取数据产生网络开销。应观察Spark UI的Locality Level若大量为NODE_LOCAL或RACK_LOCAL考虑调整spark.locality.wait参数。6.3 动态资源分配关闭导致资源浪费生产环境强烈建议开启动态资源分配配合External Shuffle Service让空闲Executor自动释放。6.4 过早使用collect()或toPandas()在调试时使用collect()或toPandas()会将所有数据拉到Driver大数据集下必定OOM。应使用take()或采样。6.5 忽视小文件问题写入数据时如果分区数过多或每个分区数据量太小会产生大量小文件给HDFS NameNode带来压力。应使用coalesce()控制输出文件数。七、完整实战案例本案例将从一个慢速的Spark作业开始逐步应用各种调优技巧最终使其性能提升5倍以上。我们模拟一个常见的电商场景计算每个商品的销售总额并关联商品类别。7.1 初始代码未优化# optimization_before.py # 未优化的版本全量读取、Shuffle大、数据倾斜未处理frompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportcol,sumasspark_sumimporttime sparkSparkSession.builder \.appName(OptimizationBefore)\.master(yarn)\.getOrCreate()# 生成模拟数据1亿条销售记录其中某个商品IDP_9999占比30%倾斜# 实际生产应从Hive读取这里简化data[]foriinrange(100_000_000):product_idfP_{i%100000}ifi%30:# 约33%的倾斜product_idP_9999amountrandom.randint(1,1000)data.append((product_id,amount))dfspark.createDataFrame(data,[product_id,amount])# 简单的groupBy聚合starttime.time()resultdf.groupBy(product_id).agg(spark_sum(amount).alias(total_amount))result.count()# 触发行动print(f执行耗时:{time.time()-start}秒)问题单次Shuffle数据倾斜严重未使用AQE未设置合适的shuffle分区数。7.2 优化后代码# optimization_after.py # 优化版本AQE、广播变量、加盐打散、参数调优frompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportcol,sumasspark_sum,concat,lit,rand,split,exprfrompyspark.sql.typesimport*importtime# 创建SparkSession时配置大量优化参数sparkSparkSession.builder \.appName(OptimizationAfter)\.master(yarn)\.config(spark.sql.shuffle.partitions,200)\.config(spark.sql.adaptive.enabled,true)\.config(spark.sql.adaptive.coalescePartitions.enabled,true)\.config(spark.sql.adaptive.skewJoin.enabled,true)\.config(spark.sql.autoBroadcastJoinThreshold,10485760)\.config(spark.serializer,org.apache.spark.serializer.KryoSerializer)\.config(spark.sql.adaptive.skewedJoinThreshold,10485760)\.config(spark.dynamicAllocation.enabled,true)\.config(spark.dynamicAllocation.minExecutors,5)\.config(spark.dynamicAllocation.maxExecutors,50)\.getOrCreate()# 生成相同的数据略与之前一致# 方法1直接依靠AQE处理倾斜Spark 3.x自动处理starttime.time()resultdf.groupBy(product_id).agg(spark_sum(amount).alias(total_amount))result.count()print(f仅开启AQE耗时:{time.time()-start}秒)# 方法2手动加盐打散更彻底print(\n使用加盐打散...)# 给product_id加随机盐0-9salted_dfdf.withColumn(salted_key,concat(col(product_id),lit(_),(rand()*10).cast(int)))# 第一次聚合partialsalted_df.groupBy(salted_key).agg(spark_sum(amount).alias(partial_sum))# 去掉盐二次聚合finalpartial.withColumn(product_id,split(col(salted_key),_)[0])\.groupBy(product_id).agg(spark_sum(partial_sum).alias(total_amount))final.count()print(f加盐打散耗时:{time.time()-start}秒)# 验证结果一致性spark.stop()7.3 SQL调优示例假设我们有订单表orders和用户表users需要统计每个用户的订单总额。# 未优化SQLsql SELECT u.user_id, SUM(o.amount) as total FROM orders o JOIN users u ON o.user_id u.user_id GROUP BY u.user_id dfspark.sql(sql)# 优化后启用AQE使用广播Join如果users表小于10MB# 自动识别小表并广播无需改写7.4 监控与诊断通过Spark UI观察Stages页面查看长尾Task确认数据倾斜。Storage页面查看缓存命中率。SQL页面查看物理计划确认是否使用了BroadcastJoin、分区裁剪等。Executors页面查看GC时间、Shuffle读写量判断内存是否不足。八、代码逐行解析8.1 资源配置参数.config(spark.sql.shuffle.partitions,200)设置Shuffle分区数避免默认200对于小数据量过多或对于大数据量过少。8.2 AQE关键参数.config(spark.sql.adaptive.enabled,true).config(spark.sql.adaptive.coalescePartitions.enabled,true).config(spark.sql.adaptive.skewJoin.enabled,true)AQE会在运行时动态优化合并小分区、处理倾斜Join、动态切换Join策略。8.3 加盐打散实现salted_dfdf.withColumn(salted_key,concat(col(product_id),lit(_),(rand()*10).cast(int)))partialsalted_df.groupBy(salted_key).agg(spark_sum(amount))finalpartial.withColumn(product_id,split(salted_key,_)[0]).groupBy(product_id).agg(spark_sum(partial_sum))通过随机盐将倾斜key分散到10个分区先局部聚合再去盐全局聚合。注意盐的粒度需根据倾斜程度调整。九、业务场景落地应用9.1 场景一大表Join大表优化当两张大表Join时无法广播。常见优化使用分桶bucket预先存储数据避免Shuffle。使用相同的分区器让数据在物理上对齐。使用Bloom Filter预过滤。9.2 场景二多维度聚合中的倾斜使用rollup或cube时某些维度组合可能产生大量数据。可使用spark.sql.adaptive.skewJoin.enabled自动处理。9.3 场景三流计算中的数据倾斜Structured Streaming中窗口聚合也可能倾斜。可以使用加盐或增加分区。9.4 场景四机器学习特征工程中的倾斜在One-Hot编码或TF-IDF中高频词可能导致倾斜。使用局部聚合或采样。十、常见报错排查10.1Container killed by YARN for exceeding memory limits原因Executor实际内存超过申请值包括堆外内存。解决增加spark.executor.memoryOverhead或减少spark.executor.memory。10.2FetchFailedException原因Shuffle数据拉取失败可能是节点故障或网络问题。解决增加重试次数spark.shuffle.io.maxRetries开启外部shuffle服务。10.3OutOfMemoryError: GC overhead limit exceeded原因JVM GC时间过长通常因为内存中对象过多如缓存了大量数据。解决减少缓存数据量使用MEMORY_AND_DISK或增加内存。10.4org.apache.spark.shuffle.MetadataFetchFailedException原因Shuffle阶段某个Map任务输出丢失。解决检查Executor日志排查磁盘故障增加任务重试次数。十一、本节课知识点总结调优参数速查表类别参数推荐值说明资源spark.executor.memory8g-16g堆内存资源spark.executor.cores3-5每个Executor的核数资源spark.dynamicAllocation.enabledtrue动态资源分配并行度spark.sql.shuffle.partitions200-500SQL shuffle分区数并行度spark.default.parallelism2-3倍总核数RDD默认并行度内存spark.memory.fraction0.6-0.8统一内存占比序列化spark.serializerKryoSerializer性能提升SQLspark.sql.adaptive.enabledtrue开启AQESQLspark.sql.autoBroadcastJoinThreshold10m-50m广播阈值Shufflespark.shuffle.file.buffer64k缓冲区大小Shufflespark.reducer.maxSizeInFlight96m拉取数据大小数据倾斜治理决策树是否Join引起倾斜 ├─ 是 → 小表是否10MB │ ├─ 是 → 广播Join │ └─ 否 → 拆分倾斜key单独处理 或 加盐打散Join需特殊处理 └─ 否聚合引起 → 加盐两阶段聚合性能优化流程基准测试运行作业记录耗时和资源消耗。监控分析通过Spark UI定位瓶颈长尾Task、大Shuffle、高GC。参数调优调整资源配置、并行度、内存、序列化。代码优化减少Shuffle、使用内置函数、避免UDF、优化Join。倾斜处理识别倾斜key选择合适方案。验证迭代对比优化前后效果。十二、课后思考作业作业一理论理解题解释Spark统一内存管理模型中执行内存和存储内存的借用机制以及可能导致的相互淘汰问题。AQE是如何动态合并Shuffle分区的它如何判断一个分区是否需要合并请列举至少三种数据倾斜的场景并分别说明适合的解决方案。作业二代码实践题编写一个PySpark程序生成一个严重倾斜的数据集一个key占80%分别使用加盐两阶段聚合和AQE自动处理Spark 3.x对比执行时间和Shuffle数据量。使用explain查看一个Join的执行计划判断是否使用了BroadcastJoin。如果没有如何强制广播模拟一个Shuffle Spill的场景内存不足导致溢写磁盘调整spark.sql.shuffle.partitions和spark.memory.fraction观察溢写量的变化。作业三场景应用题某日志分析平台每天处理TB级日志核心计算是SELECT device_id, COUNT(*) FROM logs GROUP BY device_id。目前任务运行在100节点集群每个节点16核64G发现Shuffle阶段严重倾斜少数设备ID如空字符串占50%数据。请给出完整的优化方案包括如何识别倾斜的设备ID加盐打散的具体实现写出代码。资源配置参数建议executor数量、内存、分区数等。如何验证优化效果作业四拓展研究深入研究Spark SQL的AQE源码了解CoalesceShufflePartitions和OptimizeSkewedJoin的实现原理。学习使用Spark的Stage和Task级别的Metrics如执行时间、GC时间、Shuffle读写量编写一个自动化分析脚本识别慢任务并给出优化建议。对比Presto/Trino与Spark SQL在处理数据倾斜上的不同策略写一份调研报告。提交方式本次作业要求提交代码、运行日志截图、Spark UI截图以及调优前后对比数据。扩展阅读Spark官方文档Tuning Guide《Spark性能优化实战》- 阿里巴巴技术团队Spark源码sql/core/src/main/scala/org/apache/spark/sql/execution/adaptive通过本节课的学习你已经掌握了PySpark性能调优的全套方法论从资源到代码从参数到架构。学完这节课你应该能够从容应对生产环境中的各种性能问题。下一节课我们将进行综合项目实战结合所有知识完成一个完整的大数据分析项目。我们下节课见《20节课 PySpark 从入门到精通》系列课程导航去订阅 感谢您耐心阅读到这里 如果本文对您有所启发欢迎 点赞 收藏 分享给更多需要的伙伴。️ 期待在评论区看到您的想法, 共同进步。 关注我持续获取更多干货内容 我们下篇文章见