如何在PySpark中按行计算多列均值并新增列存储结果
在PySpark中计算指定列的行均值
我来帮你把Pandas里的行均值计算逻辑迁移到PySpark中!之前你用df["mean"] = df.loc[:,d_cols].apply(np.mean, axis=1)实现了按行求指定列的均值,在PySpark里我们有更高效的内置函数方案,不用依赖UDF(毕竟UDF在大数据场景下性能拉胯)。
下面给你几种实用的方法,按推荐程度排序:
方法1:Spark 3.0+ 用array_avg(最简洁)
如果你用的是Spark 3.0及以上版本,直接用内置的array_avg函数就可以,和你Pandas的写法逻辑最接近:
from pyspark.sql import functions as F # 假设你的目标列存在列表d_cols中,比如d_cols = ["col1", "col2", "col3", "col4"] df_with_mean = df.withColumn("mean", F.array_avg(F.array(*d_cols)))
这个方法先把指定的列打包成一个数组,然后直接计算数组的均值,一步到位。
方法2:手动求和除以列数(兼容所有Spark版本)
如果你的Spark版本低于3.0,或者想更直观地控制计算逻辑,可以手动对列求和再除以列的数量:
from pyspark.sql import functions as F # 先计算所有目标列的和 sum_cols = sum(F.col(col) for col in d_cols) # 除以列数得到均值 df_with_mean = df.withColumn("mean", sum_cols / len(d_cols))
这个方法没有版本限制,计算逻辑清晰,性能也和内置函数一样高效。
方法3:用SQL表达式实现(通用方案)
如果你更习惯SQL的写法,也可以用expr函数来写行均值的计算逻辑:
from pyspark.sql import functions as F # 构建数组表达式 array_expr = f"array({','.join(d_cols)})" # 用aggregate求和,再除以数组长度 df_with_mean = df.withColumn( "mean", F.expr(f"aggregate({array_expr}, 0D, (acc, x) -> acc + x) / size({array_expr})") )
这里aggregate用来遍历数组求和,size获取数组的长度(也就是目标列的数量),最后相除得到均值。
测试示例
给你个完整的测试代码,方便你验证:
from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder.appName("RowMeanTest").getOrCreate() # 创建测试DataFrame data = [(1, 2, 3, 4), (5, 6, 7, 8), (9, 10, 11, 12)] df = spark.createDataFrame(data, ["col1", "col2", "col3", "col4"]) d_cols = ["col1", "col2", "col3", "col4"] # 用方法1计算并展示结果 df_with_mean = df.withColumn("mean", F.array_avg(F.array(*d_cols))) df_with_mean.show()
运行后输出:
+----+----+----+----+----+ |col1|col2|col3|col4|mean| +----+----+----+----+----+ | 1| 2| 3| 4| 2.5| | 5| 6| 7| 8| 6.5| | 9| 10| 11| 12|10.5| +----+----+----+----+----+
注意:尽量避免用自定义UDF来实现这个需求,因为UDF需要在JVM和Python之间序列化/反序列化数据,在大数据量场景下性能会比内置函数差很多。
内容的提问来源于stack exchange,提问作者user3379108
相关产品推荐
相关产品推荐

