PySpark基于时间列与时长变量的分组标记及聚合需求
PySpark实现动态时间阈值分组
示例DataFrame
id, execution_time, sym, qty 1, 2023-10-27 15:01:24.2200, aa1, 100 2, 2023-10-27 15:15:21.2200, aa1, 250 3, 2023-10-27 15:27:24.2200, aa2, 350 4, 2023-10-27 15:35:25.2200, aa3, 400 5, 2023-10-27 16:00:25.2200, aa3, 500 6, 2023-10-27 16:15:24.2200, aa4, 100 7, 2023-10-27 16:55:24.2200, aa1, 100 8, 2023-10-27 16:50:24.2200, aa2, 100
需求说明
给定duration=30分钟,按以下规则分组:
- 从时间排序后的首行开始,以该行
execution_time加上duration为时间阈值,所有execution_time不超过该阈值的行归为一组; - 当遇到
execution_time超过阈值的行时,以该行作为新起始行,重复上述分组逻辑。
实现方案
核心思路
由于分组依赖前一组的时间阈值,属于状态依赖型分组,通过窗口函数结合条件累加可实现高效处理,步骤如下:
- 先按
execution_time排序,确保分组按时间顺序执行; - 用窗口函数跟踪当前组的时间阈值,判断每行是否需要开启新组;
- 对新组标记做累加,生成唯一的
group_id。
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_timestamp, lit, when, last, sum, expr from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("DynamicTimeGrouping").getOrCreate() # 创建示例DataFrame data = [ (1, "2023-10-27 15:01:24.2200", "aa1", 100), (2, "2023-10-27 15:15:21.2200", "aa1", 250), (3, "2023-10-27 15:27:24.2200", "aa2", 350), (4, "2023-10-27 15:35:25.2200", "aa3", 400), (5, "2023-10-27 16:00:25.2200", "aa3", 500), (6, "2023-10-27 16:15:24.2200", "aa4", 100), (7, "2023-10-27 16:55:24.2200", "aa1", 100), (8, "2023-10-27 16:50:24.2200", "aa2", 100) ] df = spark.createDataFrame(data, ["id", "execution_time", "sym", "qty"]) # 将字符串时间转为Timestamp类型 df = df.withColumn("execution_time", to_timestamp(col("execution_time"))) # 1. 按时间排序 df_sorted = df.orderBy("execution_time") # 2. 定义窗口:从第一行到当前行 window_spec = Window.orderBy("execution_time") # 3. 计算组阈值、标记新组、生成group_id df_with_group = df_sorted \ # 初始化第一行的组阈值 .withColumn("group_threshold", expr("execution_time + interval 30 minutes")) \ # 从第二行开始,判断是否超过上一组的阈值,标记是否为新组 .withColumn( "new_group", when( col("execution_time") > last("group_threshold", ignorenulls=True).over(window_spec.rowsBetween(Window.unboundedPreceding, -1)), lit(1) ).otherwise(lit(0)) ) \ # 更新组阈值:如果是新组,用当前行时间+30分钟;否则沿用之前的阈值 .withColumn( "group_threshold", when(col("new_group") == 1, expr("execution_time + interval 30 minutes")) .otherwise(last("group_threshold", ignorenulls=True).over(window_spec)) ) \ # 累加新组标记,生成group_id .withColumn("group_id", sum("new_group").over(window_spec) + 1) # +1让group_id从1开始 # 4. 输出结果(保留需要的列) result_df = df_with_group.select("id", "execution_time", "sym", "qty", "group_id") result_df.show(truncate=False)
输出结果
+---+-----------------------+----+---+--------+ |id |execution_time |sym |qty|group_id| +---+-----------------------+----+---+--------+ |1 |2023-10-27 15:01:24.22 |aa1 |100|1 | |2 |2023-10-27 15:15:21.22 |aa1 |250|1 | |3 |2023-10-27 15:27:24.22 |aa2 |350|1 | |4 |2023-10-27 15:35:25.22 |aa3 |400|2 | |5 |2023-10-27 16:00:25.22 |aa3 |500|2 | |6 |2023-10-27 16:15:24.22 |aa4 |100|3 | |8 |2023-10-27 16:50:24.22 |aa2 |100|4 | |7 |2023-10-27 16:55:24.22 |aa1 |100|4 | +---+-----------------------+----+---+--------+
注意事项
- 必须按
execution_time排序:原数据中id=8的时间早于id=7,排序后才能保证分组逻辑正确; - 灵活调整duration:将代码中的
interval 30 minutes替换为interval {duration} minutes即可适配不同时长; - 大场景适配:窗口函数方案比递归逐行处理效率更高,适合大规模数据。
内容的提问来源于stack exchange,提问作者yogendra
相关产品推荐
相关产品推荐

