You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 01:15:32