PySpark中如何按分组实现数据前向填充与后向填充
PySpark 分组排序后实现前向填充、后向填充
核心逻辑
基于PySpark原生窗口函数实现,无需自定义UDF,性能最优,核心逻辑如下:
- 以分组字段作为窗口分区键,排序字段作为窗口内排序规则
- 前向填充:取窗口内当前行之前的最近非缺失值填充当前缺失位
- 后向填充:取窗口内当前行之后的最近非缺失值填充当前缺失位
- 注意:PySpark窗口函数的空值忽略逻辑仅对
null生效,如果数据中缺失值为NaN类型,需要先转换为null再执行填充。
测试数据构建
from pyspark.sql import functions as F from pyspark.sql import Window # 构建原始测试数据 df = spark.createDataFrame([ ('a', 1.0, 1.0), ('b', 1.0, 2.0), ('a', 2.0, float("nan")), ('b', 2.0, float("nan")), ('a', 3.0, 3.0), ('b', 3.0, 4.0)], ["id", "order", "values"]) # 将NaN值转换为null,适配窗口函数空值忽略规则 df_with_null = df.withColumn("values", F.when(F.isnan("values"), None).otherwise(F.col("values")))
原始数据预览:
+---+-----+------+ | id|order|values| +---+-----+------+ | a| 1.0| 1.0| | b| 1.0| 2.0| | a| 2.0| NaN| | b| 2.0| NaN| | a| 3.0| 3.0| | b| 3.0| 4.0| +---+-----+------+
前向填充(Forward Fill)实现
前向填充的窗口范围为分区内第一行到当前行,使用last函数取窗口内最近的非空值填充:
# 定义前向填充窗口 ffill_window = Window.partitionBy("id").orderBy("order").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 执行前向填充 ffill_df = df_with_null.withColumn("values", F.last("values", ignorenulls=True).over(ffill_window)) ffill_df.show()
执行结果和预期一致:
+---+-----+------+ | id|order|values| +---+-----+------+ | a| 1.0| 1.0| | b| 1.0| 2.0| | a| 2.0| 1.0| | b| 2.0| 2.0| | a| 3.0| 3.0| | b| 3.0| 4.0| +---+-----+------+
后向填充(Backward Fill)实现
后向填充的窗口范围为当前行到分区内最后一行,使用first函数取窗口内最近的非空值填充:
# 定义后向填充窗口 bfill_window = Window.partitionBy("id").orderBy("order").rowsBetween(Window.currentRow, Window.unboundedFollowing) # 执行后向填充 bfill_df = df_with_null.withColumn("values", F.first("values", ignorenulls=True).over(bfill_window)) bfill_df.show()
执行结果和预期一致:
+---+-----+------+ | id|order|values| +---+-----+------+ | a| 1.0| 1.0| | b| 1.0| 2.0| | a| 2.0| 3.0| | b| 2.0| 4.0| | a| 3.0| 3.0| | b| 3.0| 4.0| +---+-----+------+
内容的提问来源于stack exchange,提问作者Mykola Zotko
相关产品推荐
相关产品推荐

