You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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版本中被官方修复。

解决方案

  1. 优先升级版本:直接升级到PySpark 3.5.0及以上版本,彻底解决该性能问题。
  2. 临时兼容方案(无法升级时):
    改用数组函数将目标列转为数组,再通过数组聚合函数计算总和,避免生成大量嵌套加法表达式:
    # 方案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
    )
    
    这种方式能绕过3.4.x版本的优化缺陷,大幅提升计算性能。

内容的提问来源于stack exchange,提问作者krishna

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 12:45:20