Spark3执行及写入速度远慢于Spark2的性能调优问询
Spark 2 与 Spark 3 性能差异调优方案
问题场景
- 三个Spark DataFrame(
util、dept、cntry)从带WHERE条件的Hive表读取,均包含customer_id列 - 循环合并为
combined_spark_df:- 每个DataFrame选择
customer_id及其他列 - 通过
unionAll追加到初始为None的combined_spark_df
- 每个DataFrame选择
- 在
combined_spark_df上执行mapPartition返回8个Pandas DataFrame,后续处理:- 结果RDD转为Spark DataFrame(
results_rdd.toDF(schema)) - 执行
results_df.persist()得到8个Spark DataFrame
- 结果RDD转为Spark DataFrame(
- 后续操作的性能差异:
- 某DataFrame的count+窗口过滤:Spark2/Spark3均耗时约13分钟
- 另一DataFrame的count操作:Spark2耗时2.71秒,Spark3耗时8分钟
- 缓存所有DataFrame后写入Hive表:Spark2写入速度快(缓存生效),Spark3中
mapPartition重复执行,写入耗时极长
- 总耗时:Spark2约35分钟,Spark3约155分钟;已尝试调优:内存+磁盘持久化、关闭AQE、缓存后count、启用PyArrow,效果不佳
后续调优方向
1. 优化DataFrame合并逻辑
- 替换
unionAll为unionByName(若列顺序可能不一致),或确保三个源DataFrame列顺序完全一致。Spark3对Schema一致性校验更严格,不一致会触发额外转换开销。 - 避免循环
union,直接使用df1.union(df2).union(df3)或spark.createDataFrame(sc.union([df1.rdd, df2.rdd, df3.rdd]), schema),循环追加会生成嵌套执行计划,Spark3优化器难以有效处理。
2. 改进mapPartition与持久化策略
- 不在
mapPartition中返回多个Pandas DataFrame,改为每个分区输出单一结构化数据(如Tuple/自定义Row),再拆分出8个DataFrame。Spark3跨语言数据传递逻辑有变化,多DataFrame输出会增加序列化开销,易触发重复计算。 - 持久化时指定
StorageLevel.MEMORY_AND_DISK_SER序列化存储,减少内存占用,适配Spark3内存管理模型;执行persist()后立即调用count()触发缓存落地,避免惰性求值延迟缓存。
3. 替换df.rdd.isEmpty()为高效API
- Spark3.0+支持
df.isEmpty(),直接利用DataFrame统计信息判断,无需触发全量计算;若必须用RDD方式,确保DataFrame已完全缓存后再调用df.rdd.isEmpty()。
4. 排查执行计划与分区策略
- 用
explain()对比Spark2/Spark3的count操作执行计划,检查是否存在额外Shuffle、数据倾斜。手动设置spark.sql.shuffle.partitions与Spark2一致(默认200),避免Spark3分区优化逻辑导致的性能损耗。 - 调整
mapPartition后DataFrame的分区数,根据集群CPU核数设置为核数的2-3倍,避免并行度过高/过低。
5. 优化Hive读写逻辑
- 用
explain()验证Hive表读取时WHERE条件是否正确下推到存储层(如Parquet/ORC分区过滤),确保分区过滤生效。 - 写入Hive表时,开启
spark.sql.hive.convertMetastoreParquet=true,指定ORC/Parquet格式并启用压缩;用spark.sql("INSERT INTO ...")替代DataFrame的writeAPI,减少客户端与集群交互开销。
6. 优化Python-JVM交互配置
- 确认PyArrow版本与Spark3兼容(Spark3.2+建议PyArrow 6.0+),设置
spark.sql.execution.arrow.pyspark.enabled=true和spark.sql.execution.arrow.pyspark.fallback.enabled=false,避免低效序列化 fallback。 - 将
mapPartition替换为Spark SQL的pandas_udf(如groupBy().applyInPandas()),利用Spark3的Vectorized Execution优化,减少跨语言数据传递开销。
内容的提问来源于stack exchange,提问作者Ananth Gopinath
相关产品推荐
相关产品推荐

