基于PySpark按日历时间15分钟间隔拆分数据
按日历对齐的时间间隔拆分数据
需要将数据按对齐日历的固定时间间隔(如15/30分钟)拆分,而非基于每条数据起始时间的偏移块。
示例输入数据
ID | rh_start_time | rh_end_time | total_duration 5421833835 | 31-12-2023 13:26:53 | 31-12-2023 13:27:03 | 10 5421833961 | 31-12-2023 13:23:50 | 31-12-2023 13:39:10 | 360
期望输出(30分钟日历间隔)
ID | rh_start_time | rh_end_time | total_duration | Interval Start 5421833835 | 31-12-2023 13:26:53 | 31-12-2023 13:27:03 | 10 | 31-12-2023 13:00:00 5421833961 | 31-12-2023 13:23:50 | 31-12-2023 13:39:10 | 360 | 31-12-2023 13:00:00 5421833961 | 31-12-2023 13:23:50 | 31-12-2023 13:39:10 | 360 | 31-12-2023 13:30:00
问题点
直接用explode+sequence会基于每条数据的起始时间生成偏移间隔(如从13:26:53开始加30分钟得到13:56:53),无法对齐到日历的整点/半点等固定间隔。
解决方案
核心思路是先将每条数据的起始、结束时间对齐到日历的固定间隔,再生成该区间内的间隔序列并展开。
PySpark 实现代码
from pyspark.sql import functions as F # 1. 先将字符串时间转换为Timestamp类型(如果原始数据是字符串格式) df = df.withColumn("rh_start_time", F.to_timestamp("rh_start_time", "dd-MM-yyyy HH:mm:ss")) df = df.withColumn("rh_end_time", F.to_timestamp("rh_end_time", "dd-MM-yyyy HH:mm:ss")) # 2. 计算对齐到日历间隔的起始/结束点(这里用30分钟间隔,需要15分钟的话替换成"15 minutes") df_aligned = df.withColumn( "aligned_start", F.date_trunc("30 minutes", F.col("rh_start_time")) ).withColumn( "aligned_end", # 处理结束时间刚好在间隔边界的情况,确保包含最后一个覆盖的间隔 F.when( (F.minute(F.col("rh_end_time")) % 30 == 0) & (F.second(F.col("rh_end_time")) == 0), F.date_trunc("30 minutes", F.col("rh_end_time")) ).otherwise( F.date_trunc("30 minutes", F.col("rh_end_time")) + F.expr("interval 30 minutes") ) ) # 3. 生成间隔序列并展开 result_df = df_aligned.withColumn( "Interval Start", F.explode(F.sequence(F.col("aligned_start"), F.col("aligned_end"), F.expr("interval 30 minutes"))) ).drop("aligned_start", "aligned_end") # 4. 转换回原时间格式(可选,根据需求调整) result_df = result_df.withColumn("rh_start_time", F.date_format("rh_start_time", "dd-MM-yyyy HH:mm:ss")) \ .withColumn("rh_end_time", F.date_format("rh_end_time", "dd-MM-yyyy HH:mm:ss")) \ .withColumn("Interval Start", F.date_format("Interval Start", "dd-MM-yyyy HH:mm:ss")) # 查看结果 result_df.show(truncate=False)
代码说明
date_trunc("30 minutes", time):将时间截断到最近的30分钟日历间隔(如13:26:53→13:00:00,13:39:10→13:30:00),这是实现日历对齐的关键。- 处理
aligned_end的逻辑:确保结束时间所在的间隔被包含,比如13:39:10对应的对齐结束点是14:00:00,这样序列会生成13:00:00和13:30:00两个间隔,覆盖整个时间段。 - 如果需要15分钟间隔,只需把代码中所有的
"30 minutes"替换为"15 minutes"即可。
内容的提问来源于stack exchange,提问作者its_niks
相关产品推荐
相关产品推荐

