基于时间间隔为PySpark DataFrame分配分组ID
PySpark 按时间间隔分组生成唯一GROUP_ID解决方案
处理逻辑
- 将字符串格式的时间列转换为Timestamp类型,便于时间差计算
- 按
UID分区、Time排序,计算每行与上一行的时间间隔(小时) - 标记时间间隔超过指定阈值(如3小时)的行作为新组起点
- 对每个
UID内的标记值累加,得到组内局部分组ID - 通过全局排序生成唯一的
GROUP_ID
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("TimeGrouping").getOrCreate() # 创建示例DataFrame data = [ (1, "10/1/2016 7:25:52 AM"), (1, "10/1/2016 8:53:38 AM"), (1, "10/1/2016 11:18:50 AM"), (1, "10/1/2016 3:19:32 PM"), (2, "10/1/2016 10:25:36 AM"), (2, "10/1/2016 10:28:08 AM"), (3, "10/1/2016 10:57:41 AM"), (3, "10/1/2016 8:57:10 PM") ] df = spark.createDataFrame(data, ["UID", "Time"]) # 1. 转换时间列为Timestamp类型 df = df.withColumn("Time", F.to_timestamp("Time", "M/d/yyyy h:mm:ss a")) # 2. 定义分区窗口,计算时间差 uid_window = Window.partitionBy("UID").orderBy("Time") df = df.withColumn( "time_diff_hours", F.coalesce((F.unix_timestamp("Time") - F.unix_timestamp(F.lag("Time").over(uid_window))) / 3600, F.lit(0)) ) # 3. 标记新组(间隔超过3小时则为新组) N = 3 df = df.withColumn("is_new_group", F.when(F.col("time_diff_hours") > N, 1).otherwise(0)) # 4. 生成组内局部ID df = df.withColumn("local_group_id", F.sum("is_new_group").over(uid_window) + 1) # 5. 生成全局唯一GROUP_ID global_window = Window.orderBy("UID", "local_group_id") df = df.withColumn("GROUP_ID", F.dense_rank().over(global_window)) # 展示结果 result_df = df.select( "UID", F.date_format("Time", "M/d/yyyy h:mm:ss a").alias("Time"), "GROUP_ID" ) result_df.show(truncate=False)
关键说明
lag("Time").over(uid_window):获取同一用户上一条操作的时间,用于计算时间间隔unix_timestamp:将时间转换为秒数,方便计算时间差并转换为小时单位- 累加
is_new_group:实现同一用户内连续操作的分组,时间间隔超过阈值时触发新组 dense_rank():确保全局范围内每个组的GROUP_ID唯一,不会出现断层
内容的提问来源于stack exchange,提问作者mht
相关产品推荐
相关产品推荐

