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

如何不使用partitionBy向前填充PySpark DataFrame部分列?

解决PySpark DataFrame缺失值向下填充问题

需求说明

需要将DataFrame各列的null值,用该列最近的非null值向下填充,直到遇到下一个非null值为止。

实现方案

利用PySpark的窗口函数结合last()函数即可实现,last()的ignoreNulls=True参数会忽略空值,取窗口内最后一个有效非空值。

完整代码示例

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import last

# 初始化Spark会话
spark = SparkSession.builder.appName("FillMissingValues").getOrCreate()

# 构建原始DataFrame
raw_data = [
    (1, None, None),
    (None, "A", 99),
    (None, None, None),
    (None, None, None),
    (None, "B", 100),
    (None, None, None),
    (None, None, None),
    (None, "C", 101),
    (None, None, None),
    (None, None, None)
]
df = spark.createDataFrame(raw_data, ["column1", "column2", "column3"])

# 定义窗口:按原始行序,范围从第一行到当前行
window_spec = Window.orderBy(spark.sql("monotonically_increasing_id")).rowsBetween(Window.unboundedPreceding, Window.currentRow)

# 对每一列应用填充逻辑
filled_df = df.withColumn("column1", last("column1", ignoreNulls=True).over(window_spec)) \
              .withColumn("column2", last("column2", ignoreNulls=True).over(window_spec)) \
              .withColumn("column3", last("column3", ignoreNulls=True).over(window_spec))

# 打印结果
filled_df.show()

关键代码解释

  • monotonically_increasing_id():生成唯一递增ID,确保窗口按原始数据的行顺序排序,避免行序混乱。
  • Window.unboundedPreceding, Window.currentRow:设置窗口范围为从数据起始行到当前行,保证每次取到的是当前行之前最近的非空值。
  • last(col, ignoreNulls=True):在窗口内忽略空值,提取最后一个有效非空值,实现向下填充效果。

运行代码后即可得到你期望的填充结果。

内容的提问来源于stack exchange,提问作者royalewithcheese

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 17:03:17