Spark DataFrame筛选符合时间间隔的重复诊疗记录及代码优化
优化Spark DataFrame患者医疗操作筛选逻辑
需求说明
我有一个包含患者ID(Pat_ID)和医疗操作日期(Date)的DataFrame,需要筛选出满足以下条件的记录:
- 患者至少接受过两次医疗操作
- 存在两次操作的间隔在30至365天之间
- 仅保留符合时间范围条件的首次操作对应的患者ID和日期
原DataFrame示例
| Pat_ID | Date |
|---|---|
| A | 1-march-18 |
| A | 15-march-18 |
| B | 1-april-19 |
| B | 4-april-19 |
| B | 7-april-19 |
| B | 3-june-19 |
筛选后目标DataFrame
| Pat_ID | Date |
|---|---|
| B | 7-april-19 |
低效实现及问题
我之前尝试了下面的代码标记目标日期,虽然能运行但效率极低——循环生成365列会严重拖慢计算速度,且占用大量内存:
w=Window.partitionBy("Pat_ID").orderBy(col("date")) for i in range(1, 366): df = df.withColumn(f"daysbetween_{i}", when ((datediff((F.lead(F.col('dx_date'), i).over(w)), "dx_date").between(30, 365)),1).otherwise(0))
高效优化方案
方案思路
避免循环生成冗余列,利用Spark的窗口函数和集合操作一次性完成判断:
- 统一日期格式为日期类型,确保计算精度
- 按患者分区,收集该患者所有操作日期到数组
- 对每个日期,检查数组中是否存在后续操作在当前日期+30至+365天范围内
- 筛选出符合条件的记录后,对每个患者取最早的目标日期
代码实现
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 转换字符串日期为日期类型(格式串可根据实际数据调整) df = df.withColumn("Date", F.to_date("Date", "d-MMM-yy")) # 2. 按患者分区,收集所有操作日期到数组 patient_window = Window.partitionBy("Pat_ID") df_with_all_dates = df.withColumn("all_operation_dates", F.collect_list("Date").over(patient_window)) # 3. 判断当前日期是否有后续操作在30-365天范围内 df_with_valid_flag = df_with_all_dates.withColumn( "has_valid_followup", F.exists( "all_operation_dates", lambda date: F.datediff(date, F.col("Date")).between(30, 365) ) ).filter(F.col("has_valid_followup")) # 4. 对每个患者取最早的符合条件的操作日期 rank_window = Window.partitionBy("Pat_ID").orderBy("Date") final_result = df_with_valid_flag.withColumn( "rank", F.row_number().over(rank_window) ).filter(F.col("rank") == 1).drop("all_operation_dates", "has_valid_followup", "rank") # 查看结果 final_result.show()
方案优势
- 避免循环生成大量冗余列,大幅降低内存占用
- 利用Spark内置的分布式集合操作,充分发挥集群计算能力
- 逻辑清晰,代码可维护性更强
内容的提问来源于stack exchange,提问作者Sarah Naeger
相关产品推荐
相关产品推荐

