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

无需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|
+---+-------+-------+-----+----+---+

错误伪代码问题说明

  1. 聚合函数误用:原代码直接使用F.sum是分组聚合逻辑,但这里需要的是行内多列的求和计算,正确方式是将多个列的判断结果作为参数传入F.sum(用*解包列表)。
  2. 语法错误:原代码将生成器表达式直接放在withColumn参数中,会导致语法异常,需先构建完整的表达式再传入。
  3. 除零未处理:未考虑count为0时的除法异常,需用F.coalesce将除法得到的null转换为0。

内容的提问来源于stack exchange,提问作者Zakaria Hamane

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 12:25:58