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

如何在PySpark DataFrame中计算与上次值变更的时间差

解决方案

要实现需求中的time_alive列,核心思路是先定位每个分区内最近的has_changed=1的时间点,再用当前行的time减去该时间点;如果分区内无has_changed=1的记录,则用分区首行的time计算差值。具体实现步骤如下:

完整代码实现

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

spark = SparkSession.builder.getOrCreate()

# 定义数据与Schema
data = [("A", 1, 0), ("A", 2, 1), ("A", 4, 0), ("A", 7, 1), ("B", 2, 1), ("B", 5, 0)]
schema = ["object", "time", "has_changed"]
df = spark.createDataFrame(data, schema)

# 定义窗口:按object分区,time升序排列
window_part = Window.partitionBy("object").orderBy("time")

# 步骤1:标记has_changed=1的行的time,其他行为null
df = df.withColumn("change_time", f.when(f.col("has_changed") == 1, f.col("time")))

# 步骤2:向前填充最近的change_time(忽略null值)
df = df.withColumn("last_change_time", f.last("change_time", ignoreNulls=True).over(window_part))

# 步骤3:处理分区内无has_changed=1的情况,用分区首行的time替代
window_first = Window.partitionBy("object").orderBy("time").rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
df = df.withColumn("first_time", f.first("time").over(window_first))
df = df.withColumn("base_time", f.coalesce(f.col("last_change_time"), f.col("first_time")))

# 步骤4:计算time_alive
df = df.withColumn("time_alive", f.col("time") - f.col("base_time"))

# 展示结果
df.select("object", "time", "has_changed", "time_alive").show()

代码说明

  1. 标记变化点时间:用when函数筛选出has_changed=1的行,记录其time到change_time列,其他行留空。
  2. 向前填充最近变化点:使用last函数结合ignoreNulls=True,在分区内按时间排序后,为每行填充最近的非空change_time,得到last_change_time。
  3. 处理无变化点的情况:用first函数获取分区首行的time,再通过coalesce函数优先取last_change_time,若为空则取首行时间作为计算基准base_time。
  4. 计算最终差值:用当前行的time减去base_time,得到time_alive。

执行上述代码后,输出结果与预期完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 03:45:03