Spark DataFrame如何删除存在空值的AccountID对应的全部记录
实现方案
你可以通过分组标记+过滤的逻辑实现需求,以下是不同场景下的实现代码:
PySpark 实现
方法1:分组过滤后关联(适合大数据量场景,shuffle效率更高)
from pyspark.sql import functions as F # 1. 统计每个AccountID下是否存在任意空值记录 valid_accounts = df.groupBy("AccountID") \ .agg( # 逐行判断是否存在空值,再取分组内最大值,为True代表该AccountID存在空值记录 F.max(F.greatest(*[F.col(col).isNull() for col in df.columns])).alias("has_any_null") ) \ .filter(F.col("has_any_null") == False) \ .select("AccountID") # 2. 关联原表得到所有符合要求的行 result_df = df.join(valid_accounts, on="AccountID", how="inner") result_df.show()
方法2:窗口函数实现(无需关联,代码更简洁)
from pyspark.sql import functions as F from pyspark.sql.window import Window # 按AccountID分区,标记该分区是否存在空值记录 window_acc = Window.partitionBy("AccountID") df = df.withColumn( "has_any_null", F.max(F.greatest(*[F.col(col).isNull() for col in df.columns])).over(window_acc) ) # 过滤仅保留无空值的AccountID对应行,删除辅助标记列 result_df = df.filter(F.col("has_any_null") == False).drop("has_any_null") result_df.show()
Spark SQL 实现
WITH account_null_flag AS ( SELECT AccountID, -- 标记该AccountID是否存在任意空值行 MAX(CASE WHEN Name IS NULL OR Price IS NULL THEN 1 ELSE 0 END) AS has_any_null FROM 你的表名 GROUP BY AccountID ) SELECT t.* FROM 你的表名 t INNER JOIN account_null_flag f ON t.AccountID = f.AccountID WHERE f.has_any_null = 0
以上实现运行后都会得到你期望的结果,仅保留AccountID为12的所有行。
内容的提问来源于stack exchange,提问作者Priyanka
相关产品推荐
相关产品推荐

