如何在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()
代码说明
- 标记变化点时间:用
when函数筛选出has_changed=1的行,记录其time到change_time列,其他行留空。 - 向前填充最近变化点:使用
last函数结合ignoreNulls=True,在分区内按时间排序后,为每行填充最近的非空change_time,得到last_change_time。 - 处理无变化点的情况:用
first函数获取分区首行的time,再通过coalesce函数优先取last_change_time,若为空则取首行时间作为计算基准base_time。 - 计算最终差值:用当前行的
time减去base_time,得到time_alive。
执行上述代码后,输出结果与预期完全一致。
内容的提问来源于stack exchange,提问作者Dani
相关产品推荐
相关产品推荐

