如何不使用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
相关产品推荐
相关产品推荐

