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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 00:20:06