无需UDF,使用PySpark实现指定列的行内计数、求和与平均值计算
纯PySpark内置函数实现多列有效数值统计与计算
需求说明
实现一个函数,接收PySpark DataFrame和列名列表,完成以下操作:
- 新增
count列:统计每行中值不等于-1的列的数量 - 新增
sum列:计算每行中值不等于-1的数值之和 - 新增
avg列:计算有效数值的平均值,无有效数值时返回0
正确实现代码
from pyspark.sql import DataFrame import pyspark.sql.functions as F def average_columns(df: DataFrame, values: list) -> DataFrame: # 构造count列表达式:每行有效列的数量 count_expr = F.sum(*[F.when(F.col(col) != -1, 1).otherwise(0) for col in values]) # 构造sum列表达式:每行有效数值的总和 sum_expr = F.sum(*[F.when(F.col(col) != -1, F.col(col)).otherwise(0) for col in values]) # 构造avg列表达式:处理除零情况,无有效数值时返回0.0 avg_expr = F.coalesce(sum_expr / count_expr, F.lit(0.0)) return df.withColumn("count", count_expr) \ .withColumn("sum", sum_expr) \ .withColumn("avg", avg_expr)
示例使用
示例输入数据构造
data = [(1, 2.0, 3.0), (2, 1.0, 3.0), (3, -1.0, 3.0), (4, -1.0, -1.0)] df = spark.createDataFrame(data, ["id", "column1", "column2"])
调用函数并查看结果
result_df = average_columns(df, ["column1", "column2"]) result_df.show()
输出结果
+---+-------+-------+-----+----+---+ | id|column1|column2|count| sum|avg| +---+-------+-------+-----+----+---+ | 1| 2.0| 3.0| 2| 5.0|2.5| | 2| 1.0| 3.0| 2| 4.0|2.0| | 3| -1.0| 3.0| 1| 3.0|3.0| | 4| -1.0| -1.0| 0| 0.0|0.0| +---+-------+-------+-----+----+---+
错误伪代码问题说明
- 聚合函数误用:原代码直接使用
F.sum是分组聚合逻辑,但这里需要的是行内多列的求和计算,正确方式是将多个列的判断结果作为参数传入F.sum(用*解包列表)。 - 语法错误:原代码将生成器表达式直接放在
withColumn参数中,会导致语法异常,需先构建完整的表达式再传入。 - 除零未处理:未考虑
count为0时的除法异常,需用F.coalesce将除法得到的null转换为0。
内容的提问来源于stack exchange,提问作者Zakaria Hamane
相关产品推荐
相关产品推荐

