PySpark计算多列均值:忽略空值且保留原列空值的实现
解决PySpark多列均值计算忽略Null值的问题
原代码直接将所有列相加后除以列数,只要任意列存在Null值,相加结果就会变成Null,最终均值列也为Null,无法满足保留原列Null同时计算有效均值的需求。要解决这个问题,需要计算每行中非Null列的总和,再除以该行中非Null列的数量,以下是两种可行实现方式:
方式一:使用数组过滤与聚合表达式
通过Spark内置的数组函数,先过滤掉Null值再计算均值,表达式逻辑直观:
from pyspark.sql import functions as F from pyspark.sql import types as T cols=['a','b','c','d','e','f'] # 过滤Null值后计算总和,除以非Null列的数量 find_mean = F.expr(""" aggregate( filter(array(a,b,c,d,e,f), x -> x is not null), 0D, (acc, x) -> acc + x, acc -> acc / size(filter(array(a,b,c,d,e,f), x -> x is not null)) ) """) data=[(1,2,3,4,5,None),(1,2,3,4,5,None),(5,4,3,2,1,None),(3,4,5,1,2,5)] schema=T.StructType([ T.StructField('a',T.IntegerType()), T.StructField('b',T.IntegerType()), T.StructField('c',T.IntegerType()), T.StructField('d',T.IntegerType()), T.StructField('e',T.IntegerType()), T.StructField('f',T.IntegerType()) ]) test=spark.createDataFrame(data,schema) result=test.withColumn('average',find_mean) result.display()
方式二:动态生成表达式(适合列数较多场景)
通过动态拼接SQL表达式,自动处理所有目标列,无需手动写数组元素,扩展性更强:
from pyspark.sql import functions as F from pyspark.sql import types as T cols=['a','b','c','d','e','f'] # 生成总和表达式:将Null值转为0参与求和 sum_expr = "+".join([f"coalesce({c}, 0)" for c in cols]) # 生成计数表达式:统计每行中非Null列的数量 count_expr = "+".join([f"case when {c} is not null then 1 else 0 end" for c in cols]) # 计算均值,同时处理全Null行的除以0问题 find_mean = F.expr(f"if({count_expr} = 0, null, ({sum_expr}) / ({count_expr}))") data=[(1,2,3,4,5,None),(1,2,3,4,5,None),(5,4,3,2,1,None),(3,4,5,1,2,5)] schema=T.StructType([ T.StructField('a',T.IntegerType()), T.StructField('b',T.IntegerType()), T.StructField('c',T.IntegerType()), T.StructField('d',T.IntegerType()), T.StructField('e',T.IntegerType()), T.StructField('f',T.IntegerType()) ]) test=spark.createDataFrame(data,schema) result=test.withColumn('average',find_mean) result.display()
两种方式都能得到期望结果:原列的Null值保留,均值列计算时自动忽略Null值,前3行均值为3,最后一行均值为3.3333。
内容的提问来源于stack exchange,提问作者qhrich
相关产品推荐
相关产品推荐

