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

如何在PySpark DataFrame中按顺序填充缺失值?

PySpark 实现按顺序向前填充非空group_id值

针对你需要将group_id字段的null值按row_id递增顺序,用前一个非空值填充到下一个非空值出现前的所有行的需求,这里提供高效适配大数据量的解决方案:

完整代码示例

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

# 初始化Spark会话(如果已有会话可跳过此步)
spark = SparkSession.builder.appName("FillGroupNulls").getOrCreate()

# 构造测试DataFrame(替换为你的实际数据读取逻辑即可)
test_data = [
    (1, 1),
    (2, None),
    (3, None),
    (4, None),
    (5, 5),
    (6, None),
    (7, None),
    (8, 8),
    (9, None),
    (10, None),
    (11, None),
    (12, None)
]
df = spark.createDataFrame(test_data, ["row_id", "group_id"])

# 定义窗口规则:按row_id升序,范围从第一行到当前行
fill_window = Window.orderBy("row_id").rowsBetween(Window.unboundedPreceding, Window.currentRow)

# 用last函数向前填充null,ignoreNulls=True确保只取非空的最近值
filled_df = df.withColumn(
    "group_id",
    last("group_id", ignoreNulls=True).over(fill_window)
)

# 查看结果
filled_df.show()

方案说明

  • 核心逻辑依赖Spark原生窗口函数,last(..., ignoreNulls=True)会在指定窗口范围内,自动跳过null值取最后一个非空的group_id,刚好匹配"从首次出现值生效到下一个值出现"的需求。
  • 窗口函数经过Spark优化,能够高效处理数百万行级别的数据,不会出现常规循环填充导致的性能瓶颈或内存溢出问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 07:44:55