如何用PySpark DataFrame下一行同列值替换null值?
PySpark:用下一行同列值替换Null值
问题背景
原始PySpark DataFrame如下:
| Column A | Column B |
|---|---|
| val1 | val1B |
| null | val2B |
| val2 | null |
| val3 | val3B |
需求:将DataFrame中所有列的Null值替换为正下方一行同列的值,最终结果如下:
| Column A | Column B |
|---|---|
| val1 | val1B |
| val2 | val2B |
| val2 | val3B |
| val3 | val3B |
当前进展:已统计行号并筛选出含Null值的行,但认为该步骤没必要,陷入困境。
解决方案
核心思路是利用**窗口函数lead**获取下一行的同列值,再结合coalesce函数替换Null值。由于Spark DataFrame本身是无序的,必须先添加行号来固定行顺序。
代码示例
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, coalesce, lead, row_number # 初始化SparkSession spark = SparkSession.builder.appName("ReplaceNullWithNextRow").getOrCreate() # 创建示例DataFrame data = [ ("val1", "val1B"), (None, "val2B"), ("val2", None), ("val3", "val3B") ] df = spark.createDataFrame(data, ["Column A", "Column B"]) # 步骤1:添加行号,固定行顺序 # 若无天然排序键,可改用monotonically_increasing_id()生成自增ID window_row_num = Window.orderBy("Column A") df_with_row = df.withColumn("row_num", row_number().over(window_row_num)) # 步骤2:定义窗口,用于获取下一行的值 window_lead = Window.orderBy("row_num") # 步骤3:批量处理所有列,替换Null值 processed_df = df_with_row.select( *[coalesce(col(c), lead(col(c), 1).over(window_lead)).alias(c) for c in df.columns], "row_num" ) # 步骤4:删除行号列(可选) final_df = processed_df.drop("row_num") # 展示结果 final_df.show()
代码说明
- 固定行顺序:Spark DataFrame不保证行的默认顺序,必须通过
row_number()生成行号,确保"正下方"的行是预期的目标行。 lead函数:lead(col, 1).over(window)用于获取当前行的下一行同列值,参数1表示偏移1行。coalesce函数:coalesce(col, lead_value)返回第一个非Null值,即当前列值不为Null时保留原值,为Null时用下一行的值替换。- 批量处理:通过列表推导式遍历所有列,避免逐个列手动编写逻辑,提升代码复用性。
内容的提问来源于stack exchange,提问作者newmon
相关产品推荐
相关产品推荐

