如何在PySpark DataFrame中计算每行数组的平均值?
解决Spark DataFrame数组列转平均值的问题
问题分析
你的代码存在两个关键问题:
- 函数中最后一行
df = df.withColumn('arrays', col('arrays')[0].cast('int'))完全多余,既不会影响返回的result_df,还存在未导入col的错误(需用F.col)。 - 若调用函数后未将返回值赋值给新变量,仍使用原
df,自然会看到原始数组。
修正后的代码(基于你的思路)
保留展开数组再分组聚合的逻辑,去掉无效代码,同时规范列命名:
from pyspark.sql import functions as F def average_arrays(df): # 展开数组为单个元素,避免覆盖原列 exploded_df = df.withColumn("array_element", F.explode("arrays")) # 按id分组计算平均值,并重命名结果列 result_df = exploded_df.groupBy("id").agg(F.avg("array_element").alias("arrays_avg")) return result_df # 调用函数并查看结果 result_df = average_arrays(df) result_df.show()
更高效的解法(无需Shuffle)
对于大数据量,展开数组会触发Shuffle,推荐直接用aggregate函数对每个数组内部计算平均值:
from pyspark.sql import functions as F # 方式1:使用SQL表达式 result_df = df.withColumn( "arrays_avg", F.expr("aggregate(arrays, 0.0, (acc, x) -> acc + x, acc -> acc / size(arrays))") ).drop("arrays") # 方式2:使用Spark函数API result_df = df.withColumn( "arrays_avg", F.aggregate( F.col("arrays"), F.lit(0.0), lambda acc, x: acc + x, lambda acc: acc / F.size(F.col("arrays")) ) ).drop("arrays") result_df.show()
输出结果
两种方法都会得到如下正确结果:
+---+------------------+ | id| arrays_avg| +---+------------------+ | 1| 34.55| | 2| 3.3| | 3|20.033333333333335| +---+------------------+
内容的提问来源于stack exchange,提问作者T_d
相关产品推荐
相关产品推荐

