合并两个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
相关产品推荐
相关产品推荐

