Spark新手求助:用PySpark基于上行值填充DataFrame的Col1、Col2空值
PySpark实现列空值的前向填充(Forward Fill)
针对你需要用上方行非空值填充当前列空值的需求,PySpark没有直接的ffill函数,但可以通过窗口函数实现,下面是完整的步骤和代码,适合新手理解:
1. 准备工作:导入依赖并创建测试数据
首先创建测试用的DataFrame,模拟包含空值的场景:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import last, monotonically_increasing_id, col # 初始化SparkSession spark = SparkSession.builder.appName("ForwardFillDemo").getOrCreate() # 模拟带空值的输入数据 sample_data = [ (1, 10), (None, 20), (2, None), (None, None), (3, 30) ] # 创建原始DataFrame raw_df = spark.createDataFrame(sample_data, schema=["Col1", "Col2"]) print("原始数据:") raw_df.show()
执行后原始数据输出:
+----+----+ |Col1|Col2| +----+----+ | 1| 10| |null| 20| | 2|null| |null|null| | 3| 30| +----+----+
2. 核心逻辑:用窗口函数实现前向填充
Spark的DataFrame是分布式无序的,所以首先需要给数据添加行号保证顺序,再通过窗口函数获取到当前行为止的最后一个非空值:
步骤2.1 添加行号
用monotonically_increasing_id()生成唯一递增的行号,确保数据顺序:
df_with_row = raw_df.withColumn("row_id", monotonically_increasing_id())
步骤2.2 定义窗口规则
窗口范围设置为从第一行到当前行,按行号排序:
# 定义窗口:按行号排序,窗口包含从起始行到当前行的所有数据 window_spec = Window.orderBy("row_id").rowsBetween(Window.unboundedPreceding, Window.currentRow)
步骤2.3 应用前向填充
对每个需要填充的列,使用last(col, ignorenulls=True)函数,取窗口内最后一个非空值:
# 生成填充后的列 filled_df = df_with_row.select( "row_id", last("Col1", ignorenulls=True).over(window_spec).alias("Col1_filled"), last("Col2", ignorenulls=True).over(window_spec).alias("Col2_filled") ) # 合并原始列和填充列,方便对比结果 final_df = df_with_row.join(filled_df, on="row_id").drop("row_id") print("填充后数据:") final_df.show()
执行后填充结果:
+----+----+----------+----------+ |Col1|Col2|Col1_filled|Col2_filled| +----+----+----------+----------+ | 1| 10| 1| 10| |null| 20| 1| 20| | 2|null| 2| 20| |null|null| 2| 20| | 3| 30| 3| 30| +----+----+----------+----------+
3. 关键细节说明
- 行号的必要性:Spark没有天然的行顺序,必须通过
row_id或业务有序列(比如时间戳、自增ID)来保证填充顺序正确,否则结果会混乱。 last函数的参数:ignorenulls=True是核心,它会跳过窗口内的空值,只取最后一个非空值,实现前向填充效果。- 分组场景扩展:如果需要按某列分组后组内填充,只需给窗口添加
partitionBy,比如:# 按Group列分组,组内前向填充 grouped_window = Window.partitionBy("Group").orderBy("row_id").rowsBetween(Window.unboundedPreceding, Window.currentRow)
4. 优化建议
- 如果你的原始数据有业务上的有序列(比如
create_time、order_id),优先用这些列排序,比monotonically_increasing_id()更高效且符合业务逻辑。 - 若不需要保留原始列,可以直接用填充后的列替换原列,减少数据冗余。
内容的提问来源于stack exchange,提问作者Avenue Guy
相关产品推荐
相关产品推荐

