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

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超过阈值的行时,以该行作为新起始行,重复上述分组逻辑。

实现方案

核心思路

由于分组依赖前一组的时间阈值,属于状态依赖型分组,通过窗口函数结合条件累加可实现高效处理,步骤如下:

  1. 先按execution_time排序,确保分组按时间顺序执行;
  2. 用窗口函数跟踪当前组的时间阈值,判断每行是否需要开启新组;
  3. 对新组标记做累加,生成唯一的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 15:50:00