PySpark count/count_distinct统计去重值时如何将null纳入计数
PySpark 去重计数纳入Null值实现方案
问题说明
调用pyspark.sql.functions下的count()、count_distinct()做列值统计时,默认会忽略null值,和DataFrame.distinct().count()将null视为有效值的逻辑不一致。
测试复现
- 构造测试数据集
# Dataframe Creation df = spark.createDataFrame([(1,"arun","engineering",20000), (2,"manoj","finance",25000), (3,None,"accounts",None), (4,"vikram",None,None)], ["id","name","dept","salary"])
- 数据集内容
+---+------+-----------+------+ | id| name| dept|salary| +---+------+-----------+------+ | 1| arun|engineering| 20000| | 2| manoj| finance| 25000| | 3| null| accounts| null| | 4|vikram| null| null| +---+------+-----------+------+
- 默认统计代码
import pyspark.sql.functions as func df.agg(func.count("id").alias("n_ids"), func.count("name").alias("n_names"), func.count("dept"), func.count("salary"))\ .show()
- 默认输出(null值被忽略,不符合预期)
+-----+-------+-----------+-------------+ |n_ids|n_names|count(dept)|count(salary)| +-----+-------+-----------+-------------+ | 4| 3| 3| 2| +-----+-------+-----------+-------------+
该现象的原因是
count_distinct()默认仅统计非空表达式的唯一值:返回提供的表达式结果唯一且非NULL的行数,替换count()为count_distinct()会得到完全相同的输出。
- 预期输出(null值纳入去重统计)
+-----+-------+-----------+-------------+ |n_ids|n_names|count(dept)|count(salary)| +-----+-------+-----------+-------------+ | 4| 4| 4| 3| +-----+-------+-----------+-------------+
单列验证时,df.select("salary").distinct().count()可以返回正确结果3,但逐列调用distinct统计效率低,也不适合多列同时聚合的场景。
实现方法
方法1:最简洁写法(推荐,Spark 2.0+全版本支持)
利用collect_set会将null作为独立有效值纳入去重集合的特性,直接取集合长度即可得到结果,和distinct().count()逻辑完全一致,性能和原生count_distinct无差异:
import pyspark.sql.functions as func df.agg( func.size(func.collect_set("id")).alias("n_ids"), func.size(func.collect_set("name")).alias("n_names"), func.size(func.collect_set("dept")).alias("count(dept)"), func.size(func.collect_set("salary")).alias("count(salary)") ).show()
逻辑说明:
collect_set(列名)会对列值做去重,且不会过滤null,会将null作为独立元素存入结果集合size()直接取集合的元素个数,就是包含null的去重总计数
方法2:拆分计算写法(兼容极老版本Spark)
核心逻辑:包含null的去重总计数 = 非空值的去重数量 + 列中是否存在null(存在则+1,不存在则+0)
import pyspark.sql.functions as func df.agg( (func.count_distinct("id") + func.ifnull(func.count_distinct(func.when(func.col("id").isNull(), 1)), 0)).alias("n_ids"), (func.count_distinct("name") + func.ifnull(func.count_distinct(func.when(func.col("name").isNull(), 1)), 0)).alias("n_names"), (func.count_distinct("dept") + func.ifnull(func.count_distinct(func.when(func.col("dept").isNull(), 1)), 0)).alias("count(dept)"), (func.count_distinct("salary") + func.ifnull(func.count_distinct(func.when(func.col("salary").isNull(), 1)), 0)).alias("count(salary)") ).show()
方法3:哨兵值替换写法(适合值范围明确的场景)
如果可以确定列中绝对不会出现某个特定值,可以用coalesce将null替换为该哨兵值后直接做count_distinct,代码更简洁:
# 示例:明确salary列为正数,用-1作为null的替换哨兵 df.agg( func.count_distinct(func.coalesce(func.col("salary"), func.lit(-1))).alias("count(salary)") ).show()
注意:该方法必须选择业务数据中绝对不可能出现的值作为哨兵,否则会导致统计结果错误。
内容的提问来源于stack exchange,提问作者jj_coder
相关产品推荐
相关产品推荐

