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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 14:16:27