PySpark 3.4.x逐行归约性能骤降问题求助
PySpark 3.4.x多列行级求和/平均性能骤降问题解决
问题描述
在PySpark 3.4.0、3.4.1版本中,当对25列以上的字段执行行级求和(进而计算平均值)操作时,性能会出现断崖式下降,运行耗时可达4分钟左右;但在3.2.4-3.3.3以及3.5.0版本中,相同逻辑的代码运行速度极快。即使改用sum()替代reduce(add, ...)实现求和,性能问题依然存在。
复现代码
import pyspark.sql.functions as F import time from functools import reduce from operator import add from pyspark.sql import SparkSession master = "local" executor_memory = "4g" driver_memory = "4g" spark = SparkSession.builder.config("spark.master", master)\ .config("spark.executor.memory", executor_memory)\ .config("spark.driver.memory", driver_memory)\ .getOrCreate() NCOLS = 30 col_names = [f"{i}" for i in range(NCOLS)] col_vals = [F.lit(f"{i}").alias(f"{i}") for i in range(NCOLS)] df = spark.createDataFrame([('id',)], schema='id STRING') df_data = df.select(["*"] + col_vals) st=time.time() df_final = df_data.withColumn("row_avg", reduce(add, [F.col(x) for x in col_names]) / NCOLS) df_final.count() print(f"processing time is {(time.time()-st)/60}")
问题根源
这是PySpark 3.4.x版本的专属缺陷,源于表达式优化逻辑的变更:当参与行级聚合的列数超过25列时,表达式生成阶段会出现性能瓶颈,导致查询计划构建或执行耗时剧增。该问题已在3.5.0版本中被官方修复。
解决方案
- 优先升级版本:直接升级到PySpark 3.5.0及以上版本,彻底解决该性能问题。
- 临时兼容方案(无法升级时):
改用数组函数将目标列转为数组,再通过数组聚合函数计算总和,避免生成大量嵌套加法表达式:
这种方式能绕过3.4.x版本的优化缺陷,大幅提升计算性能。# 方案1:使用Spark 3.3+支持的array_sum函数 df_final = df_data.withColumn("row_avg", F.array_sum(F.array(col_names)) / NCOLS) # 方案2:使用aggregate函数兼容更早版本 df_final = df_data.withColumn( "row_avg", F.aggregate(F.array(col_names), F.lit(0.0), lambda acc, x: acc + x) / NCOLS )
内容的提问来源于stack exchange,提问作者krishna
相关产品推荐
相关产品推荐

