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,提问作者Татьяна Вахрушева
相关产品推荐
相关产品推荐

