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

PySpark如何基于非空新值更新DataFrame行并生成currentAmount列

问题原因

你写的代码有三处核心问题不符合需求:

  1. 窗口按Datetime倒序排序,违背了按时间先后取最近非空值的逻辑
  2. lag函数只能读取当前行往前偏移固定行数的值,无法处理连续多行空值的填充场景
  3. 没有单独处理首个非空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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 05:09:03