PySpark条件运行窗口:按时间差实现分组序号递增计算
PySpark 条件滑动窗口实现分组序号自增计算
逻辑说明
结合样例输入输出,实际分组规则为:按starttime升序排列数据,初始组号为1;若当前行的starttime与上一行的endtime时间差大于1秒,组号自动加1,否则保持和上一行同组。
核心实现思路是通过滑动窗口逐行判断是否触发组号递增,再对递增标记做累计求和得到最终组号,不需要复杂的嵌套排序函数。
完整实现代码
from pyspark.sql import SparkSession, Window import pyspark.sql.functions as F # 初始化Spark会话 spark = SparkSession.builder.appName("conditional_running_group").getOrCreate() # 1. 构造样例数据集 并转换字段为时间戳类型 source_data = [ ("2022-01-01 03:25:53", "2022-01-01 03:25:52"), ("2022-01-01 03:25:53", "2022-01-01 03:25:52"), ("2022-01-01 03:25:53", "2022-01-01 03:25:52"), ("2022-01-01 03:25:55", "2022-01-01 03:25:54"), ("2022-01-01 03:25:57", "2022-01-01 03:25:57") ] df = spark.createDataFrame(source_data, schema=["starttime", "endtime"]) df = df.withColumn("starttime", F.to_timestamp("starttime")) \ .withColumn("endtime", F.to_timestamp("endtime")) # 2. 定义排序窗口:按starttime升序排列,逐行处理 sort_window = Window.orderBy("starttime") # 3. 提取上一行的endtime,标记是否需要递增组号 df = df.withColumn( "prev_end", F.lag("endtime").over(sort_window) ).withColumn( "increase_flag", # 第一行没有上一行数据,标记为0不递增 F.when(F.col("prev_end").isNull(), 0) # 时间差转秒级计算,差值大于1秒标记为1,否则0 .when((F.col("starttime").cast("long") - F.col("prev_end").cast("long")) > 1, 1) .otherwise(0) ) # 4. 定义累计滑动窗口,对递增标记求和后+1得到从1开始的组号 running_window = sort_window.rowsBetween(Window.unboundedPreceding, Window.currentRow) df = df.withColumn("group", F.sum("increase_flag").over(running_window) + 1) # 5. 删除中间辅助列,输出结果 final_df = df.drop("prev_end", "increase_flag") final_df.show(truncate=False)
运行结果
执行代码后输出完全匹配预期结果:
+-------------------+-------------------+-----+ |starttime |endtime |group| +-------------------+-------------------+-----+ |2022-01-01 03:25:53|2022-01-01 03:25:52|1 | |2022-01-01 03:25:53|2022-01-01 03:25:52|1 | |2022-01-01 03:25:53|2022-01-01 03:25:52|1 | |2022-01-01 03:25:55|2022-01-01 03:25:54|2 | |2022-01-01 03:25:57|2022-01-01 03:25:57|3 | +-------------------+-------------------+-----+
注意事项
- 时间差计算时将时间戳转为long类型取秒级差值,可避免时区、毫秒级精度带来的计算误差
- 如果数据需要按业务维度分组(比如不同用户、不同设备单独计算组号),只需要在窗口定义时加上
partitionBy(对应维度字段)即可 - 该实现是标准的running window(运行时滑动窗口)逻辑,全量数据只需要做两次窗口遍历,性能远高于自定义UDF或者逐行迭代方案
内容的提问来源于stack exchange,提问作者Jon B
相关产品推荐
相关产品推荐

