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

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|
+---+-------------------+----+

去重规则:

  1. 每个id仅保留date最新的行(date无Null值);
  2. 若同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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 00:05:34