PySpark如何基于非空新值更新DataFrame行并生成currentAmount列
问题原因
你写的代码有三处核心问题不符合需求:
- 窗口按
Datetime倒序排序,违背了按时间先后取最近非空值的逻辑 lag函数只能读取当前行往前偏移固定行数的值,无法处理连续多行空值的填充场景- 没有单独处理首个非空
amountNew出现之前的行的取值逻辑
解决代码
你可以用支持忽略空值的窗口聚合函数实现填充,完整实现如下:
from pyspark.sql import Window from pyspark.sql.functions import col, first, last, when # 1. 定义按id分区、时间升序排序的基础窗口 w_ordered = Window.partitionBy("id").orderBy("Datetime") # 2. 定义覆盖id全部分区的窗口,用于取每个id第一个非空amountNew对应的amountOld值 w_full = Window.partitionBy("id") \ .orderBy("Datetime") \ .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) df_final = df_final \ # 取每个id首个非空amountNew对应的amountOld,用于填充最前面的空值行 .withColumn("head_fill_val", first(when(col("amountNew").isNotNull(), col("amountOld")), ignorenulls=True).over(w_full)) \ # 按时间顺序向前填充非空的amountNew值 .withColumn("filled_new", last(col("amountNew"), ignorenulls=True).over(w_ordered)) \ # 优先用填充后的amountNew,为空则用前面的头部填充值 .withColumn("currentAmount", when(col("filled_new").isNull(), col("head_fill_val")).otherwise(col("filled_new"))) \ # 删除中间辅助列 .drop("head_fill_val", "filled_new")
逻辑说明
last(amountNew, ignorenulls=True)会按时间顺序,取到当前行为止最近的一个非空amountNew值,自动跳过中间的空值- 对于
amountNew全为空的头部行,用每个id下第一个非空amountNew对应的amountOld值填充,正好匹配你的规则
内容的提问来源于stack exchange,提问作者Mobin Ranjbar
相关产品推荐
相关产品推荐

