如何在Databricks中基于3天时间窗为DataFrame分组生成序列ID?
解决方案:按分组生成3天时间窗口序列ID
以下分别提供Pandas和PySpark两种实现方案,适配Databricks环境,解决你遇到的连续3天被分配不同序列ID的问题。
Pandas 实现
适用于中小数据量场景,逻辑清晰易调试:
import pandas as pd # 确保日期列格式为datetime df['date'] = pd.to_datetime(df['date']) # 按分组字段+日期排序,保证窗口计算的顺序性 df = df.sort_values(['service', 'phone_number', 'date']) # 计算当前行与同组上一行的日期差(首次行差值设为0) df['date_diff'] = df.groupby(['service', 'phone_number'])['date'].diff().dt.days.fillna(0) # 标记新窗口:当日期差超过2天时,视为新窗口起点 df['new_window'] = (df['date_diff'] > 2).astype(int) # 累加新窗口标记生成序列ID(每组从1开始计数) df['seq'] = df.groupby(['service', 'phone_number'])['new_window'].cumsum() + 1 # 生成窗口说明:每个序列对应的起止日期 window_details = df.groupby(['service', 'phone_number', 'seq']).agg( window_start=('date', 'min'), window_end=('date', 'max') ).astype({'window_start': str, 'window_end': str}).reset_index() # 合并窗口说明到原数据,并格式化描述文本 df = df.merge(window_details, on=['service', 'phone_number', 'seq']) df['window_desc'] = df.apply(lambda x: f"窗口{x['seq']}: {x['window_start']} 至 {x['window_end']}", axis=1) # 清理中间临时列 df = df.drop(['date_diff', 'new_window', 'window_start', 'window_end'], axis=1)
逻辑说明
- 核心是通过日期差判断窗口连续性:同组内当前日期与前一行日期差≤2天时,属于同一3天窗口;超过2天则开启新窗口。
- 对2023-02-16、17、18这类连续日期,日期差均为1,因此会被分配同一个
seq值。
PySpark 实现
适用于大数据量场景,适配Databricks分布式计算环境:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 确保日期列格式为date类型 df = df.withColumn("date", F.to_date("date")) # 定义分组排序窗口:按service、phone_number分组,按date升序排序 group_sort_window = Window.partitionBy("service", "phone_number").orderBy("date") # 计算当前行与同组上一行的日期差(首次行差值设为0) df = df.withColumn( "date_diff", F.coalesce(F.datediff(F.col("date"), F.lag("date").over(group_sort_window)), F.lit(0)) ) # 标记新窗口:日期差>2天时标记为1,否则0 df = df.withColumn("new_window", F.when(F.col("date_diff") > 2, 1).otherwise(0)) # 累加新窗口标记生成序列ID(每组从1开始) df = df.withColumn( "seq", F.sum("new_window").over(group_sort_window.rangeBetween(Window.unboundedPreceding, 0)) + 1 ) # 生成窗口说明:每个序列对应的起止日期及格式化描述 window_details = df.groupBy("service", "phone_number", "seq")\ .agg( F.min("date").alias("window_start"), F.max("date").alias("window_end") )\ .withColumn( "window_desc", F.concat( F.lit("窗口"), F.col("seq"), F.lit(": "), F.date_format("window_start", "yyyy-MM-dd"), F.lit(" 至 "), F.date_format("window_end", "yyyy-MM-dd") ) ) # 合并窗口说明到原数据 df = df.join(window_details, on=["service", "phone_number", "seq"], how="left") # 清理临时列 df = df.drop("date_diff", "new_window", "window_start", "window_end") # 查看结果 df.display() # Databricks环境推荐用display替代show
逻辑说明
- 用PySpark窗口函数
lag获取同组前一行日期,datediff计算天数差,通过sum over窗口累加新窗口标记生成序列ID。 - 完全适配分布式计算,处理亿级数据无压力,且能正确识别连续3天为同一窗口。
问题根源分析
你之前的代码(@mozway)可能误用了rank()/row_number()这类基于排序的排名函数,这类函数仅按日期顺序生成序号,未考虑时间窗口的连续性,导致连续3天被分配不同ID。本文方案通过日期差判断窗口边界,从根本上解决了这个问题。
内容的提问来源于stack exchange,提问作者deps
相关产品推荐
相关产品推荐

