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

基于时间戳时长合并行的PySpark复杂转换实现方案咨询

问题描述

我有一个包含time(时间戳)列和count(整数)列的Delta Lake表,需要对DataFrame的行进行分组合并,规则如下:

  • 按动态2天间隔分组:每组以第一行的时间戳为基准,当后续行的时间戳与该基准的差值超过172800秒(即2天)时,该行成为新组的起始基准
  • 分组后对每组的count列求和

原始数据示例

time, count
2019-02-18 11:03:55, 500
2019-02-18 11:06:18, 30
2019-02-18 11:07:58, 20
2019-02-18 11:07:58, 12
2019-02-18 11:08:38, 8
2019-02-18 11:10:29, 2
2019-02-20 11:09:12, 25
2019-02-20 11:10:10, 10
2019-04-02 10:10:10, 1
2019-04-05 10:10:10, 2
2019-04-09 10:10:09, 4
2019-04-11 10:10:30, 6
2019-04-13 10:10:10, 3
2019-04-16 10:10:10, 5
2019-04-19 10:10:10, 7
2019-04-21 10:10:10, 8

预期结果

time => count_sum
2019-02-18 11:03:55 => 572
2019-02-20 11:09:12 => 35
2019-04-02 10:10:10 => 1
2019-04-05 10:10:10 => 2
2019-04-09 10:10:09 => 4
2019-04-11 10:10:30 => 9
2019-04-16 10:10:10 => 5
2019-04-19 10:10:10 => 15
解决方案

这个需求属于动态时间窗口分组,无法用常规固定窗口函数实现,需要通过标记分组边界的方式构建分组ID,以下是基于PySpark的实现步骤:

步骤1:读取Delta表并排序

首先读取Delta表,必须按time列排序,保证分组逻辑按时间顺序执行:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化SparkSession(未初始化时执行)
spark = SparkSession.builder.appName("DynamicTimeGrouping").getOrCreate()

# 读取Delta表
df = spark.read.format("delta").load("/path/to/your/delta/table")

# 按时间排序,确保分组逻辑正确
df_sorted = df.orderBy("time")

步骤2:标记分组基准

使用窗口函数对比当前行与上一组基准时间的差值,动态更新分组基准:

# 定义全局排序窗口
window_spec = Window.orderBy("time")

# 计算每组的基准时间:当前行与上一组基准差超过2天时,以当前时间作为新基准
df_with_group = df_sorted.withColumn(
    "prev_group_base",
    F.lag("group_base", 1).over(window_spec)
).withColumn(
    "group_base",
    F.when(
        # 第一行或时间差超过2天,更新基准
        F.col("prev_group_base").isNull() | 
        (F.unix_timestamp("time") - F.unix_timestamp("prev_group_base") > 172800),
        F.col("time")
    ).otherwise(F.col("prev_group_base"))
)

步骤3:按基准分组求和

以group_base为分组键,对count列求和:

result_df = df_with_group.groupBy("group_base").agg(
    F.sum("count").alias("count_sum")
).orderBy("group_base")

# 展示结果
result_df.show(truncate=False)

验证结果

运行后得到的结果与预期一致:

group_basecount_sum
2019-02-18 11:03:55572
2019-02-20 11:09:1235
2019-04-02 10:10:101
2019-04-05 10:10:102
2019-04-09 10:10:094
2019-04-11 10:10:309
2019-04-16 10:10:105
2019-04-19 10:10:1015

注意事项

  • 必须确保数据按time排序,否则分组逻辑会出错
  • 数据量极大时,可先按月份等粗粒度时间分区,再在分区内执行上述逻辑,提升性能
  • 若time为Timestamp类型,直接用F.unix_timestamp转换为秒数计算差值即可

内容的提问来源于stack exchange,提问作者Ganesha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 00:59:56