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

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)

代码分步解释

  1. 数据准备:把原始数据转换成Spark DataFrame,并将timestamp字段转换成时间类型,方便后续时间计算。
  2. 窗口函数匹配后续0记录:通过last窗口函数向前填充,让每条indicator=1的记录都能拿到后面最近的indicator=0的时间和id,这样我们就知道每个1记录需要拆分到哪个时间点。
  3. 过滤目标记录:只保留indicator=1且存在后续0记录的行,避免处理无后续0的无效数据。
  4. 计算时间序列结束点:
    • 特殊情况:如果下一个0的时间在当前1记录的5分钟区间内,直接把结束时间设为当前时间+5分钟
    • 正常情况:计算不超过下一个0时间的最大5分钟整点,比如12:37会被处理成12:35
  5. 生成并拆分时间序列:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 18:57:38