PySpark:含Null值列用max/min去重丢失行的问题及优化问询
Spark DataFrame去重逻辑优化方案
一、问题回顾
原始DataFrame df 结构如下:
+---+-------------------+----+ | id| date|site| +---+-------------------+----+ |100|2020-03-24 00:00:00| a| |100|2019-08-30 00:00:00| a| |100|2020-03-24 00:00:00| b| |101|2019-12-20 00:00:00|NULL| |101|2019-12-20 00:00:00| a| |102|2019-04-14 00:00:00|NULL| |103|2019-09-28 00:00:00| c| +---+-------------------+----+
去重规则:
- 每个
id仅保留date最新的行(date无Null值); - 若同
id同最新date有多条记录,优先保留site非Null的行,无则保留Null行。
你已通过窗口函数筛选出各id的最新date记录(得到df2),但后续用max('site')筛选时丢失了仅含Null site的id=102记录,随后通过左反连接+union补全了数据:
df_left_anti = df.join(df3, df['id'] == df3['id'], 'left_anti') df_all = df3.union(df_left_anti).orderBy('id')
二、现有方法的正确性与效率评估
正确性
该方法完全正确:左反连接能精准提取出df中未被df3包含的id(即仅存Null site的id=102),再通过union合并后能得到符合规则的完整结果。
效率
该方法存在性能损耗:需要执行一次左反连接(涉及shuffle操作)和一次union+orderBy(再次触发shuffle),当数据量较大时,多次shuffle会显著增加计算资源消耗和执行时间,不是最优方案。
三、兼容Null值的优化方案
方案1:给Null值设置兜底值适配max/min函数
通过coalesce将site的Null值替换为一个不会影响排序结果的兜底值(比如空字符串),让max/min函数能识别并保留这类记录:
from pyspark.sql import functions as f w = Window.partitionBy('id') df3 = df2.withColumn( 'maxSite', f.max(f.coalesce(f.col('site'), f.lit(''))).over(w) ).where( f.coalesce(f.col('site'), f.lit('')) == f.col('maxSite') ).drop('maxSite')
此方法能保留id=102的记录,同时满足优先选非Null site的规则。
方案2:单窗口排序一步完成(最优)
直接通过一次窗口排序,同时满足“取最新date”和“优先非Null site”的规则,无需分多步操作:
from pyspark.sql import functions as f # 窗口规则:按id分组,先按date降序(取最新),再按site是否为Null升序(非Null在前) w = Window.partitionBy('id').orderBy( f.col('date').desc(), f.col('site').isNull().asc() ) df_final = df.withColumn('row_num', f.row_number().over(w)) \ .where(f.col('row_num') == 1) \ .drop('row_num')
逻辑说明:
date.desc确保每个id只保留最新的日期记录;site.isNull().asc将非Null的site排在前面(因为site.isNull()返回True为1、False为0,升序排序时0在前);row_number()取每组的第一条记录,完美匹配需求。
该方案仅需一次窗口操作,无额外shuffle,性能最优且代码简洁。
内容的提问来源于stack exchange,提问作者TDG
相关产品推荐
相关产品推荐

