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

如何在PySpark中对分组后的向量/数组求均值?

Spark分组计算数组列均值的正确实现方法

当Spark DataFrame包含数组类型列,按类别分组后直接调用avg()函数会报错——因为Spark内置聚合函数默认不支持数组/向量类型的直接聚合。以下是几种高效可行的实现方式:

方法1:利用高阶函数聚合(推荐,通用型)

通过collect_list收集每组的数组集合,再用aggregate和transform按位置累加后求均值,适配任意固定长度的数组:

from pyspark.sql import functions as F

df1.groupBy("c").agg(
    F.aggregate(
        F.collect_list("v1"),
        # 初始值:创建与目标数组长度一致的全0数组
        F.array(*[F.lit(0.0) for _ in range(2)]),
        # 累加阶段:将当前数组与累加器按位置相加
        lambda acc, arr: F.transform(
            F.arrays_zip(acc, arr),
            lambda x: x["acc"] + x["arr"]
        ),
        # 计算均值阶段:累加结果除以组内元素数量
        lambda acc: F.transform(acc, lambda x: x / F.count("v1"))
    ).alias("v1")
).show()

方法2:拆分数组元素后聚合(适合短数组)

如果数组长度固定且较短,可以先将数组元素拆分为单独列,用普通avg聚合后再合并为数组:

from pyspark.sql import functions as F

# 拆分数组元素为独立列
df_split = df1.select(
    "c",
    F.col("v1")[0].alias("v1_0"),
    F.col("v1")[1].alias("v1_1")
)

# 分组求均值后合并回数组
df_split.groupBy("c").agg(
    F.avg("v1_0").alias("avg_0"),
    F.avg("v1_1").alias("avg_1")
).select(
    "c",
    F.array("avg_0", "avg_1").alias("v1")
).show()

方法3:自定义UDF(不推荐,性能较低)

若需要快速验证逻辑,可以用Python UDF结合Numpy计算,但UDF会触发数据序列化到Python端,性能远低于内置函数:

from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, DoubleType
import numpy as np

@F.udf(returnType=ArrayType(DoubleType()))
def array_avg(arr_list):
    return np.mean(arr_list, axis=0).tolist()

df1.groupBy("c").agg(
    array_avg(F.collect_list("v1")).alias("v1")
).show()

输出结果

三种方法均会得到预期结果:

+---+----------+
|  c|        v1|
+---+----------+
|  1|[2.5, 3.5]|
+---+----------+

内容的提问来源于stack exchange,提问作者Ícaro Lorran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 15:42:17