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

合并两个DataFrame并按AusID保留最新时间记录的实现方案

解决方案提示

针对你百万级DataFrame的去重需求,这里有一套高效的处理流程,分三步就能解决:

1. 统一Timestamp格式

因为两个DataFrame的Timestamp格式不同,第一步必须把它们转换成相同的Timestamp类型,否则后续无法正确比较时间先后。用Spark的to_timestamp函数指定各自的格式即可:

# 示例:根据你的实际时间格式替换字符串
from pyspark.sql.functions import to_timestamp

# 假设df1的Timestamp格式是"yyyy-MM-dd HH:mm:ss",df2是"dd/MM/yyyy HH:mm"
df1 = df1.withColumn("Timestamp", to_timestamp("Timestamp", "yyyy-MM-dd HH:mm:ss"))
df2 = df2.withColumn("Timestamp", to_timestamp("Timestamp", "dd/MM/yyyy HH:mm"))

注意:如果是带毫秒的格式,记得加上SSS;不确定格式的话,可以先通过df.select("Timestamp").show(5)查看样本数据

2. 合并两个DataFrame

优先用unionByName合并(比union更安全,能自动对齐列顺序,避免因列顺序不同导致的数据错位):

combined_df = df1.unionByName(df2)

如果两个DataFrame的列名、列顺序完全一致,用union也可以,但unionByName更稳妥,适合列较多的场景。

3. 按AusID去重,保留最新Timestamp的记录

这里用窗口函数是最高效的方案(适配百万级数据规模,性能远优于groupBy后再join的方式):

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number

# 定义窗口规则:按AusID分组,每组内按Timestamp降序排序
window_spec = Window.partitionBy("AusID").orderBy(combined_df["Timestamp"].desc())

# 给每组内的记录加行号,最新的记录行号为1
ranked_df = combined_df.withColumn("row_num", row_number().over(window_spec))

# 筛选出每组行号为1的记录,就是每个AusID的最新数据
final_df = ranked_df.filter(ranked_df["row_num"] == 1).drop("row_num")

为什么你的原方法不适用?

你之前用的df1.union(df2).except(df1.intersect(df2))是基于所有列完全匹配来判断重复的,而你的需求是仅按AusID判断重复、保留时间最新的记录——这两个逻辑完全不匹配,所以这个方法无法满足你的需求,窗口函数才是正确的思路。

另外针对百万级数据的小建议:

  • 如果是Spark环境,尽量给Timestamp或AusID列设置分区,能大幅提升窗口函数的执行效率
  • 可以根据集群资源调整Spark的executor内存参数,避免出现内存溢出问题

内容的提问来源于stack exchange,提问作者Waqar Ahmed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 03:59:05