在PySpark中保留Patient_id对应最大time_stamp记录并删除重复
解决Databricks中按Patient_id去重并保留最新时间戳记录的问题
直接用dropDuplicates()没法实现按时间戳筛选保留最新记录的需求,因为它只能基于指定列去重,无法自定义保留规则。这里提供两种实用方法:
方法一:窗口函数(推荐,逻辑直观)
利用row_number()窗口函数给每个Patient_id分组内的记录按时间戳降序标行号,取行号为1的记录就是每个患者最新的那条。
from pyspark.sql import Window from pyspark.sql.functions import row_number, desc # 定义窗口规则:按Patient_id分区,time_stamp降序排序 window_spec = Window.partitionBy("Patient_id").orderBy(desc("time_stamp")) # 添加行号、筛选最新记录、删除行号列 deduped_df = df.withColumn("row_num", row_number().over(window_spec)) \ .filter("row_num = 1") \ .drop("row_num")
如果同一个Patient_id存在多条相同最大时间戳的记录,这种方法会只保留其中一条;如果需要保留所有同最大时间戳的记录,可以把row_number()换成rank()。
方法二:分组聚合+关联
先分组找出每个Patient_id对应的最大时间戳,再通过关联原表获取完整记录。
from pyspark.sql.functions import max # 分组计算每个患者的最新时间戳 max_time_df = df.groupBy("Patient_id").agg(max("time_stamp").alias("max_time")) # 关联原表筛选出对应记录,清理冗余列 deduped_df = df.join(max_time_df, (df.Patient_id == max_time_df.Patient_id) & (df.time_stamp == max_time_df.max_time), "inner") \ .drop(max_time_df.Patient_id, max_time_df.max_time)
这种方法会保留所有与最大时间戳匹配的记录,适合允许同患者同时间多条记录的场景。
注意事项
- 确保
time_stamp列是时间类型(如TimestampType),如果是字符串格式,先转换类型:df = df.withColumn("time_stamp", df["time_stamp"].cast("timestamp")) - 两种方法的性能差异不大,窗口函数更易理解和维护,分组关联在超大规模数据集下可能有轻微性能优势,按需选择即可。
内容的提问来源于stack exchange,提问作者Darkmaster
相关产品推荐
相关产品推荐

