如何在PySpark/Python中按姓名分组计算前两条记录的日期差
按姓名分组计算日期间隔实现方案
下面分别给出Pandas(Python原生数据处理库)和PySpark两种环境下的实现代码,逻辑完全匹配需求规则。
方案1:Pandas实现
核心逻辑:
- 先将日期字段转为标准datetime格式,避免日期计算错误
- 按姓名字段分组,组内按日期升序排序
- 每个分组判断记录条数:仅1条记录时间隔固定为0;2条及以上记录时,取排序后前2条日期做差得到间隔天数,将该值赋值给组内所有行
- 整理输出格式得到最终结果
完整代码:
import pandas as pd # 加载/构造输入数据 df = pd.DataFrame({ "Name": ["Nancy", "Rictk", "Francky", "Nancy", "Nancy", "Francky"], "Date": ["2021-08-14", "2021-08-15", "2021-08-16", "2021-08-18", "2022-02-07", "2021-12-06"] }) # 日期字段类型转换 df["Date"] = pd.to_datetime(df["Date"]) def calc_group_diff(group_df): # 组内按日期升序排列 sorted_dates = group_df.sort_values("Date", ascending=True)["Date"].values if len(sorted_dates) < 2: group_df["Day"] = 0 else: # 计算前两条记录的天数差 day_diff = (sorted_dates[1] - sorted_dates[0]).astype('timedelta64[D]').astype(int) group_df["Day"] = day_diff return group_df # 分组应用计算逻辑 result = df.groupby("Name", group_keys=False).apply(calc_group_diff) # 日期转成字符串格式方便输出 result["Date"] = result["Date"].dt.strftime("%Y-%m-%d") print(result)
方案2:PySpark实现
核心逻辑:
- 日期字段转为Spark的Date类型
- 通过窗口函数按姓名分区、日期排序,给组内记录打上排序序号
- 分组聚合时,提取每个姓名排序后第1、2条记录的日期,用
datediff函数计算间隔,记录数不足2条时直接返回0 - 将计算得到的间隔值关联回原明细表,得到所有行带Day字段的结果
完整代码:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark会话 spark = SparkSession.builder.appName("date_diff_calc").getOrCreate() # 加载/构造输入数据 source_data = [ ("Nancy", "2021-08-14"), ("Rictk", "2021-08-15"), ("Francky", "2021-08-16"), ("Nancy", "2021-08-18"), ("Nancy", "2022-02-07"), ("Francky", "2021-12-06") ] df = spark.createDataFrame(source_data, schema=["Name", "Date"]) df = df.withColumn("Date", F.to_date(F.col("Date"), "yyyy-MM-dd")) # 定义分区排序窗口 sort_window = Window.partitionBy("Name").orderBy(F.col("Date").asc()) # 计算每个姓名对应的间隔天数 diff_mapping = df.withColumn("rn", F.row_number().over(sort_window)) \ .groupBy("Name") \ .agg( F.when(F.count("*") < 2, 0) .otherwise( F.datediff( F.max(F.when(F.col("rn") == 2, F.col("Date"))), F.max(F.when(F.col("rn") == 1, F.col("Date"))) ) ).alias("Day") ) # 关联回原表得到最终结果 result = df.join(diff_mapping, on="Name", how="left") \ .withColumn("Date", F.date_format(F.col("Date"), "yyyy-MM-dd")) result.show()
注:你给出的示例输出中Francky的日期间隔为6属于笔误,按原始输入数据Francky排序后的前两条日期为2021-08-16和2021-12-06,实际间隔为112天,上述代码完全遵循你描述的计算规则。
内容的提问来源于stack exchange,提问作者Shubham Korade
相关产品推荐
相关产品推荐

