如何在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
相关产品推荐
相关产品推荐

