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

使用Spark合并连续时间段的多行时间序列数据

合并Spark DataFrame中连续相同Description的时间段记录

嘿,刚好之前做过类似的时间序列合并需求,用Spark的窗口函数就能完美解决!咱们的核心思路是先给连续且Description相同的记录打上同一个分组ID,然后按分组聚合就能得到你要的结果。

1. 准备测试数据(以Python为例)

先把你提供的示例数据转换成Spark DataFrame:

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

spark = SparkSession.builder.appName("MergeContinuousRecords").getOrCreate()

data = [
    ("01:02:00", "01:05:00", "A", 1.0),
    ("01:05:00", "01:08:00", "A", 4.0),
    ("01:08:00", "01:11:00", "A", 4.3),
    ("01:11:00", "01:15:00", "B", 18.2),
    ("01:21:00", "01:55:00", "C", 0.0),
    ("01:55:00", "02:07:00", "A", 1.8)
]

df = spark.createDataFrame(data, ["Start", "End", "Description", "Value"])
df.show()

2. 核心:给连续记录打分组ID

这里需要用窗口函数做两个关键判断,来识别新分组的起点:

  • 当前记录的Description是否和上一条不同
  • 当前记录的Start是否不等于上一条记录的End(时间不连续)

只要满足其中一个条件,就标记为新分组,再通过累计求和生成唯一的分组ID:

# 先按Start排序,保证时间顺序正确,这一步至关重要
window_order = Window.orderBy("Start")
# 按Description分区,获取同类型记录的上一条End时间
window_part = Window.partitionBy("Description").orderBy("Start")

# 标记新分组的起点
df_with_flag = df.withColumn("prev_end", F.lag("End").over(window_part)) \
                 .withColumn("prev_desc", F.lag("Description").over(window_order)) \
                 .withColumn(
                     "is_new_group",
                     F.when(
                         (F.col("Description") != F.col("prev_desc")) | (F.col("Start") != F.col("prev_end")),
                         1
                     ).otherwise(0)
                 )

# 生成分组ID,对is_new_group累计求和
df_with_group = df_with_flag.withColumn(
    "group_id",
    F.sum("is_new_group").over(window_order.rangeBetween(Window.unboundedPreceding, 0))
)

df_with_group.show()

3. 按分组ID聚合得到最终结果

现在只需要按分组ID和Description聚合,取最早的Start、最晚的End,以及Value的总和:

final_df = df_with_group.groupBy("group_id", "Description") \
                        .agg(
                            F.min("Start").alias("Start"),
                            F.max("End").alias("End"),
                            F.sum("Value").alias("Value(SUM)")
                        ) \
                        .orderBy("Start") \
                        .drop("group_id")

final_df.show()

运行后就能得到你想要的结果:

+--------+--------+-----------+----------+
|   Start|     End|Description|Value(SUM)|
+--------+--------+-----------+----------+
|01:02:00|01:11:00|          A|       9.3|
|01:11:00|01:15:00|          B|      18.2|
|01:21:00|01:55:00|          C|       0.0|
|01:55:00|02:07:00|          A|       1.8|
+--------+--------+-----------+----------+

几个注意点

  • 必须先按Start排序,否则连续记录的顺序混乱会导致分组错误
  • 如果你的时间是带日期的Timestamp类型,逻辑完全一致,只需用F.to_timestamp把字符串转成时间类型即可
  • 用Scala实现的话,语法几乎相同,只是函数调用的细节略有差异

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:43:33