PySpark中如何按组将列空值替换为下一个非空值
问题原因与解决方案
问题原因
你当前使用的F.last(..., ignorenulls=True)配合默认窗口范围时,行为和预期不符,核心原因在于窗口的默认范围设置:
- 当你用
Window.partitionBy('id').orderBy('date')时,Spark窗口的默认范围是RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW,也就是只包含当前行及之前的所有行。 - 这时候
last()函数会在「当前行之前的所有行」中找最后一个非空值,自然就是前一个非空值,而不是你想要的「当前行之后的下一个非空值」。
解决方案
要获取当前行之后的下一个非空值,需要调整窗口范围为「当前行到分组末尾」,然后用first()函数(因为从当前行往后找,第一个非空值就是你要的下一个有效数据):
正确代码
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口:按id分组,date升序,窗口范围从当前行到分组最后一行 window_spec = Window.partitionBy('id').orderBy('date').rowsBetween(0, Window.unboundedFollowing) # 用first函数取当前行及之后的第一个非空last_order_date df_fixed = df.withColumn( 'last_order_date', F.first('last_order_date', ignorenulls=True).over(window_spec) )
效果验证
以id=001的行为例:
- 2021-02、2021-03的null会被替换为2021-04行的
2021-01 - 2021-05、2021-06、2021-07的null会被替换为2021-08行的
2021-04
完全符合你「用下一个非空值填充null」的需求。
补充说明
如果你的数据中存在分组末尾全为null的情况,这些null会保留(因为之后没有非空值可以填充),这是合理的预期行为。
内容的提问来源于stack exchange,提问作者Germanifold
相关产品推荐
相关产品推荐

