PySpark:基于字段值过滤Lag函数计算日期差的实现方法
PySpark按字段分组计算最近同字段行的日期差
我有一个存储字段值与日期的PySpark DataFrame,想要添加一列,用于计算每行与最近拥有相同字段值的行之间的日期差。普通连续行的差值用Lag函数就能实现,但不知道如何基于行的特定字段值进行过滤。
示例输入
1111| 23/May/2024 2222| 20/May/2024 3333| 19/May/2024 1111| 16/May/2024 4444| 12/May/2024 1111| 07/May/2024 2222| 01/May/2024
期望输出
1111| 23/May/2024| 7 2222| 20/May/2024| 19 3333| 19/May/2024| default 1111| 16/May/2024| 9 4444| 12/May/2024| default 1111| 07/May/2024| default 2222| 01/May/2024| default
解决方案
核心思路是用窗口函数+分组排序,让Lag函数只在相同字段值的组内取上一行数据,具体步骤如下:
- 转换日期格式:先把字符串日期转为PySpark可计算的
DateType类型 - 定义分组窗口:按目标字段分组,再按日期降序排序,确保同组内最新日期排在前面
- 计算日期差:用
lag获取同组上一行的日期,再用datediff计算天数差,最后把空值替换为default
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark会话 spark = SparkSession.builder.appName("GroupDateDiff").getOrCreate() # 构造示例数据 sample_data = [ ("1111", "23/May/2024"), ("2222", "20/May/2024"), ("3333", "19/May/2024"), ("1111", "16/May/2024"), ("4444", "12/May/2024"), ("1111", "07/May/2024"), ("2222", "01/May/2024") ] df = spark.createDataFrame(sample_data, ["field", "date_str"]) # 将字符串日期转为Date类型 df = df.withColumn("date", F.to_date(F.col("date_str"), "dd/MMM/yyyy")) # 定义窗口:按field分组,按date降序排列 group_window = Window.partitionBy("field").orderBy(F.col("date").desc()) # 计算同组最近日期差,处理空值 result_df = df.withColumn("prev_date", F.lag("date").over(group_window)) \ .withColumn("date_diff", F.datediff(F.col("date"), F.col("prev_date"))) \ .withColumn("date_diff", F.when(F.col("date_diff").isNull(), "default").otherwise(F.col("date_diff"))) \ .select("field", "date_str", "date_diff") # 打印结果 result_df.show(truncate=False)
代码关键点说明
partitionBy("field"):确保Lag函数只在相同字段值的分组内取上一行数据,不会跨组orderBy(F.col("date").desc()):让同组内日期最新的行排在最前面,这样Lag取到的就是当前行的最近同组日期datediff:计算两个日期的天数差,当前日期减上一行日期,结果为正整数when函数:将同组第一条数据的空差值替换为default
内容的提问来源于stack exchange,提问作者NaiveBayesian
相关产品推荐
相关产品推荐

