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

Spark DataFrame去重需求:按指定字段保留最新时间行

当然有成熟的实现方案!这种按指定字段分组、保留每组内时间最新记录的需求,在Spark处理中非常常见,我给你分享两种最常用的思路,你可以根据自己的数据规模和偏好选择:

方法一:使用窗口函数(最直观常用)

这是处理这类需求的标准做法,利用row_number()窗口函数给每组内的记录按时间排序打标,然后筛选出每组的第一条记录即可:

步骤分解:

  • 定义窗口:按CNTA_TIPODOCUMENTOS和CNTA_NRODOCUMENTO进行分组(partition)
  • 排序规则:在每个分组内,按CNTA_FECHA_FORMULARIO降序排列(这样最新的记录会排在最前面)
  • 打行号:给每条记录分配行号,每组内最新的记录行号为1
  • 过滤:只保留行号为1的记录

Python 代码示例:

from pyspark.sql import Window
from pyspark.sql.functions import row_number, desc

# 定义窗口规则
window_spec = Window.partitionBy("CNTA_TIPODOCUMENTOS", "CNTA_NRODOCUMENTO") \
                    .orderBy(desc("CNTA_FECHA_FORMULARIO"))

# 添加行号并过滤
df_unique = df.withColumn("row_num", row_number().over(window_spec)) \
              .filter("row_num == 1") \
              .drop("row_num")

Scala 代码示例:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.{row_number, desc}

val windowSpec = Window.partitionBy("CNTA_TIPODOCUMENTOS", "CNTA_NRODOCUMENTO")
                       .orderBy(desc("CNTA_FECHA_FORMULARIO"))

val dfUnique = df.withColumn("row_num", row_number().over(windowSpec))
                 .filter($"row_num" === 1)
                 .drop("row_num")

针对你举的CNTA_NRODOCUMENTO=35468731的例子,这个分组内的记录会按CNTA_FECHA_FORMULARIO降序排列,2012-08-25的那条会被标记为row_num=1,最终被保留下来。


方法二:分组取最大时间后关联

如果你的数据量特别大,或者担心窗口函数的 shuffle 开销,可以先聚合出每组的最新时间,再和原表关联筛选:

步骤分解:

  • 分组聚合:按CNTA_TIPODOCUMENTOS和CNTA_NRODOCUMENTO分组,取出每组的最大CNTA_FECHA_FORMULARIO
  • 关联筛选:将聚合结果和原DataFrame关联,只保留匹配到时间的记录(也就是每组的最新记录)

Python 代码示例:

from pyspark.sql.functions import max

# 先获取每组的最新时间
latest_dates = df.groupBy("CNTA_TIPODOCUMENTOS", "CNTA_NRODOCUMENTO") \
                 .agg(max("CNTA_FECHA_FORMULARIO").alias("latest_date"))

# 关联原表筛选最新记录
df_unique = df.join(latest_dates, 
                    (df["CNTA_TIPODOCUMENTOS"] == latest_dates["CNTA_TIPODOCUMENTOS"]) &
                    (df["CNTA_NRODOCUMENTO"] == latest_dates["CNTA_NRODOCUMENTO"]) &
                    (df["CNTA_FECHA_FORMULARIO"] == latest_dates["latest_date"]),
                    how="inner") \
              .drop(latest_dates["CNTA_TIPODOCUMENTOS"], latest_dates["CNTA_NRODOCUMENTO"], "latest_date")

Scala 代码示例:

import org.apache.spark.sql.functions.max

val latestDates = df.groupBy("CNTA_TIPODOCUMENTOS", "CNTA_NRODOCUMENTO")
                   .agg(max("CNTA_FECHA_FORMULARIO").alias("latest_date"))

val dfUnique = df.join(latestDates, 
                       df("CNTA_TIPODOCUMENTOS") === latestDates("CNTA_TIPODOCUMENTOS") &&
                       df("CNTA_NRODOCUMENTO") === latestDates("CNTA_NRODOCUMENTO") &&
                       df("CNTA_FECHA_FORMULARIO") === latestDates("latest_date"),
                       "inner")
                 .drop(latestDates("CNTA_TIPODOCUMENTOS"), latestDates("CNTA_NRODOCUMENTO"), "latest_date")

这种方法的好处是聚合的shuffle量可能比窗口函数小,但如果同一分组内有多个记录时间相同且都是最新的,会保留所有这些记录;而窗口函数的方法会只保留其中一条(如果你需要保留所有同时间的最新记录,可以把row_number()换成rank()或者dense_rank())。


小提示

  • 如果你的CNTA_FECHA_FORMULARIO字段是字符串类型,记得先转换成Timestamp类型再排序/聚合,否则可能会出现排序错误:to_timestamp(df["CNTA_FECHA_FORMULARIO"], "yyyy-MM-dd HH:mm:ss")(格式根据你的实际情况调整)
  • 两种方法各有优劣:窗口函数代码更简洁直观,适合大多数场景;分组关联在特定大数据场景下性能可能更优,你可以根据自己的测试结果选择。

内容的提问来源于stack exchange,提问作者jose rivera

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:32:38