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

PySpark DataFrame如何用后续非空值填充Null值?

在PySpark DataFrame中用后续非空值填充Null值

可以实现用后续出现的非空值替换Null值,你之前的代码无效是因为窗口范围设置错误,默认窗口只能获取当前行及之前的数值,无法拿到后续的非空值。以下是针对你需求的解决方案:

原代码问题分析

你使用的first函数配合默认窗口Window.partitionBy("product_id").orderBy("ts"),其默认窗口范围是从分区起始到当前行,只能取到当前行及之前的第一个非空值,无法获取后续的非空值,因此无法满足需求。

解决方案:用后续最近非空值填充(匹配你的预期结果)

通过标记非空行并倒序累加生成分组ID,让每个Null行与后续最近的非空行归为同一组,再用组内非空值填充所有行:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

def fill_na_prices(data):
    # 1. 标记非空price的行
    df_with_flag = data.withColumn(
        "non_null_flag",
        F.when(F.col("price").isNotNull(), 1).otherwise(0)
    )
    
    # 2. 倒序累加标记值,生成分组ID:让Null行与后续最近非空行同组
    window_rev = Window.partitionBy("product_id").orderBy(F.desc("ts")).rowsBetween(Window.unboundedPreceding, 0)
    df_with_group = df_with_flag.withColumn(
        "group_id",
        F.sum("non_null_flag").over(window_rev)
    )
    
    # 3. 按分组取非空price,填充组内所有行
    window_group = Window.partitionBy("product_id", "group_id")
    df_filled = df_with_group.withColumn(
        "price",
        F.first("price", ignorenulls=True).over(window_group)
    ).drop("non_null_flag", "group_id")
    
    return df_filled

# 测试执行
filled_df = fill_na_prices(data)
filled_df.orderBy("product_id", "ts").show()

执行后得到预期结果:

+----------+----------+-----+
|product_id|        ts|price|
+----------+----------+-----+
|         1|2024-05-01|  109|
|         1|2024-05-02|  109|
|         1|2024-05-03|  120|
|         2|2024-05-01|  115|
|         2|2024-05-02|  115|
|         2|2024-05-03|  115|
+----------+----------+-----+

备选方案:用后续最后一个非空值填充

如果你的需求是用后续所有非空值的最后一个填充(比如product_id=1的第一行填充120),可以直接使用last函数并设置窗口范围为当前行到分区末尾:

def fill_na_prices_last(data):
    window = Window.partitionBy("product_id").orderBy("ts").rowsBetween(0, Window.unboundedFollowing)
    return data.withColumn(
        "price",
        F.last("price", ignorenulls=True).over(window)
    )

内容的提问来源于stack exchange,提问作者Татьяна Вахрушева

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 11:17:36