使用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
相关产品推荐
相关产品推荐

