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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 08:42:20