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

如何用PySpark DataFrame下一行同列值替换null值?

PySpark:用下一行同列值替换Null值

问题背景

原始PySpark DataFrame如下:

Column AColumn B
val1val1B
nullval2B
val2null
val3val3B

需求:将DataFrame中所有列的Null值替换为正下方一行同列的值,最终结果如下:

Column AColumn B
val1val1B
val2val2B
val2val3B
val3val3B

当前进展:已统计行号并筛选出含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()

代码说明

  1. 固定行顺序:Spark DataFrame不保证行的默认顺序,必须通过row_number()生成行号,确保"正下方"的行是预期的目标行。
  2. lead函数:lead(col, 1).over(window)用于获取当前行的下一行同列值,参数1表示偏移1行。
  3. coalesce函数:coalesce(col, lead_value)返回第一个非Null值,即当前列值不为Null时保留原值,为Null时用下一行的值替换。
  4. 批量处理:通过列表推导式遍历所有列,避免逐个列手动编写逻辑,提升代码复用性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 03:05:30