PySpark实现按5分钟间隔拆分特定条件下的DataFrame行
解决PySpark按5分钟间隔拆分时间区间的需求
嘿,我来帮你搞定这个时间拆分的需求!咱们一步步来实现它,完全匹配你给出的预期输出。
需求回顾
我们需要处理indicator=1的记录,找到它后面最近的indicator=0的时间,然后按5分钟间隔拆分时间点:
- 正常情况:从
indicator=1的时间开始,每次加5分钟,直到时间点不超过下一个0的时间 - 特殊情况:如果下一个
0的时间落在当前1记录的5分钟区间内(比如13:15的下一个0在13:17),要生成该区间的所有5分钟点(包括区间结束点,比如13:20)
另外,预期输出里的id是下一条indicator=0记录的id,这个细节也需要注意。
完整实现代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("TimeIntervalSplit").getOrCreate() # 创建测试数据 data = [ (0, 128, "2019-12-03 12:00:00.0", 0), (1, 128, "2019-12-03 12:30:00.0", 1), (2, 128, "2019-12-03 12:37:00.0", 0), (3, 128, "2019-12-03 13:15:00.0", 1), (4, 128, "2019-12-03 13:17:00.0", 0) ] df = spark.createDataFrame(data, ["id", "sourceid", "timestamp", "indicator"]) df = df.withColumn("timestamp", F.to_timestamp("timestamp")) # 定义窗口:按sourceid分组,timestamp升序排列 window_spec = Window.partitionBy("sourceid").orderBy("timestamp") # 为每条记录获取后续第一个indicator=0的时间和对应的id df_with_next_zero = df \ .withColumn("next_zero_time", F.when(F.col("indicator") == 0, F.col("timestamp"))) \ .withColumn("next_zero_id", F.when(F.col("indicator") == 0, F.col("id"))) \ .withColumn("next_zero_time", F.last("next_zero_time", ignorenulls=True).over(window_spec.rowsBetween(Window.currentRow, Window.unboundedFollowing))) \ .withColumn("next_zero_id", F.last("next_zero_id", ignorenulls=True).over(window_spec.rowsBetween(Window.currentRow, Window.unboundedFollowing))) # 过滤出需要处理的indicator=1的记录,且确保存在后续的0记录 df_target = df_with_next_zero.filter(F.col("indicator") == 1).dropna(subset=["next_zero_time", "next_zero_id"]) # 计算每个记录需要生成的时间序列的结束时间 df_target = df_target.withColumn( "end_time", F.when( # 特殊情况:下一个0的时间在当前1记录的5分钟区间内 F.col("next_zero_time") <= F.col("timestamp") + F.expr("interval 5 minutes"), F.col("timestamp") + F.expr("interval 5 minutes") ).otherwise( # 正常情况:生成到不超过next_zero_time的最大5分钟整点 F.date_trunc("hour", F.col("next_zero_time")) + F.expr(f"interval ((minute(next_zero_time) // 5) * 5) minutes") ) ) # 生成时间序列并展开成多行 df_result = df_target \ .withColumn("timestamp", F.explode(F.sequence(F.col("timestamp"), F.col("end_time"), F.expr("interval 5 minutes")))) \ .select(F.col("next_zero_id").alias("id"), "sourceid", "timestamp", F.lit(1).alias("indicator")) # 查看最终结果 df_result.orderBy("id", "timestamp").show(truncate=False)
代码分步解释
- 数据准备:把原始数据转换成Spark DataFrame,并将
timestamp字段转换成时间类型,方便后续时间计算。 - 窗口函数匹配后续0记录:通过
last窗口函数向前填充,让每条indicator=1的记录都能拿到后面最近的indicator=0的时间和id,这样我们就知道每个1记录需要拆分到哪个时间点。 - 过滤目标记录:只保留
indicator=1且存在后续0记录的行,避免处理无后续0的无效数据。 - 计算时间序列结束点:
- 特殊情况:如果下一个0的时间在当前1记录的5分钟区间内,直接把结束时间设为当前时间+5分钟
- 正常情况:计算不超过下一个0时间的最大5分钟整点,比如12:37会被处理成12:35
- 生成并拆分时间序列:用
sequence函数生成从起始时间到结束时间的5分钟间隔序列,再用explode把序列拆分成多行,最后调整列名得到预期格式。
运行这段代码后,你会得到完全符合要求的输出:
+---+--------+-------------------+---------+ |id |sourceid|timestamp |indicator| +---+--------+-------------------+---------+ |1 |128 |2019-12-03 12:30:00|1 | |1 |128 |2019-12-03 12:35:00|1 | |4 |128 |2019-12-03 13:15:00|1 | |4 |128 |2019-12-03 13:20:00|1 | +---+--------+-------------------+---------+
内容的提问来源于stack exchange,提问作者nehacharya
相关产品推荐
相关产品推荐

