如何在PySpark中计算同ID内Type D与前序Type I的行间日期差?
解决方案:计算同一ID下Type D与前序Type I的日期差
PySpark 实现方案
假设你的DataFrame包含ID、Type、first_date、last_date字段,且日期字段为DateType类型。以下是具体实现步骤:
- 导入依赖模块
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, lag, datediff
- 初始化SparkSession并构造示例数据(可替换为你的实际数据源)
spark = SparkSession.builder.appName("DateDiffCalculation").getOrCreate() # 示例数据 data = [ (1, "I", "2023-01-01", "2023-01-05"), (1, "D", "2023-01-06", "2023-01-10"), (2, "I", "2023-02-01", "2023-02-03"), (2, "D", "2023-02-04", "2023-02-08"), (2, "I", "2023-02-09", "2023-02-12"), (2, "D", "2023-02-13", "2023-02-15") ] df = spark.createDataFrame(data, ["ID", "Type", "first_date", "last_date"]) # 将字符串日期转为Date类型 df = df.withColumn("first_date", col("first_date").cast("date")) df = df.withColumn("last_date", col("last_date").cast("date"))
- 定义窗口并计算日期差
# 按ID分区,按first_date排序保证行顺序正确 window_spec = Window.partitionBy("ID").orderBy("first_date") # 添加前一行的last_date字段 df = df.withColumn("prev_last_date", lag(col("last_date"), 1).over(window_spec)) # 计算Type D行与前一行Type I的天数差 result_df = df.withColumn( "days_diff", datediff(col("first_date"), col("prev_last_date")).when(col("Type") == "D", col("days_diff")).otherwise(None) ) # 查看结果 result_df.show()
说明:
lag函数用于获取同一ID分组中前一行的last_datedatediff计算两个日期的天数差- 通过
when仅保留Type D行的差值,其他行显示None
Pandas 实现方案
如果PySpark环境受限,可使用Pandas完成相同计算:
- 导入依赖并构造示例数据
import pandas as pd data = [ (1, "I", "2023-01-01", "2023-01-05"), (1, "D", "2023-01-06", "2023-01-10"), (2, "I", "2023-02-01", "2023-02-03"), (2, "D", "2023-02-04", "2023-02-08"), (2, "I", "2023-02-09", "2023-02-12"), (2, "D", "2023-02-13", "2023-02-15") ] df = pd.DataFrame(data, columns=["ID", "Type", "first_date", "last_date"]) # 转为日期类型 df["first_date"] = pd.to_datetime(df["first_date"]) df["last_date"] = pd.to_datetime(df["last_date"])
- 分组排序并计算日期差
# 按ID分组,按first_date排序 df = df.sort_values(["ID", "first_date"]).reset_index(drop=True) # 获取同一ID下前一行的last_date df["prev_last_date"] = df.groupby("ID")["last_date"].shift(1) # 计算天数差,仅保留Type D行的结果 df["days_diff"] = df.apply( lambda row: (row["first_date"] - row["prev_last_date"]).days if row["Type"] == "D" else None, axis=1 ) # 查看结果 print(df)
说明:
shift(1)用于获取分组内前一行的last_date- 通过
apply判断Type是否为D,仅计算对应行的天数差
内容的提问来源于stack exchange,提问作者StephL
相关产品推荐
相关产品推荐

