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
相关产品推荐
相关产品推荐

