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

