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

PySpark count/count_distinct统计去重值时如何将null纳入计数

PySpark 去重计数纳入Null值实现方案

问题说明

调用pyspark.sql.functions下的count()、count_distinct()做列值统计时,默认会忽略null值,和DataFrame.distinct().count()将null视为有效值的逻辑不一致。

测试复现

  1. 构造测试数据集
# 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"])
  1. 数据集内容
+---+------+-----------+------+
| id|  name|       dept|salary|
+---+------+-----------+------+
|  1|  arun|engineering| 20000|
|  2| manoj|    finance| 25000|
|  3|  null|   accounts|  null|
|  4|vikram|       null|  null|
+---+------+-----------+------+
  1. 默认统计代码
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()
  1. 默认输出(null值被忽略,不符合预期)
+-----+-------+-----------+-------------+
|n_ids|n_names|count(dept)|count(salary)|
+-----+-------+-----------+-------------+
|    4|      3|          3|            2|
+-----+-------+-----------+-------------+

该现象的原因是count_distinct()默认仅统计非空表达式的唯一值:返回提供的表达式结果唯一且非NULL的行数,替换count()为count_distinct()会得到完全相同的输出。

  1. 预期输出(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 05:09:30