PySpark:如何用分区最后值填充空值且保留非空值
PySpark 用分区最后非空值填充空值(保留中间非空值)
你当前代码的问题在于直接用last函数覆盖了所有行的取值,没有区分原列是否为空——只有原列为空时才需要用分区最后值填充,非空值必须保留原值。
解决方案
核心思路是:先拿到每个分区的最后非空值,再用coalesce函数判断,优先保留原列非空值,仅当原列为空时替换为分区最后值。
完整代码示例:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 定义分区窗口,获取每个分区的最后非空值 partition_window = Window.partitionBy(partition_column_list) df_with_last = input_df.withColumn( "partition_last_value", F.last(col(fill_column), ignorenulls=True).over(partition_window) ) # 2. 用coalesce填充空值,保留原列非空内容 filled_df = df_with_last.withColumn( f"Filled{fill_column}", F.coalesce(col(fill_column), col("partition_last_value")) ).drop("partition_last_value") # 清理中间临时列
代码解释
partition_window按分区列划分整个分区范围,last(col(fill_column), ignorenulls=True)会跳过空值,取分区内最后出现的有效数值,每个分区的所有行都会带上这个值。coalesce函数的作用是返回第一个非空参数:原列Value非空时直接用原值,为空时则用分区最后非空值填充,完全匹配你的预期输出。
原代码错误原因
你之前的窗口添加了orderBy且范围是rowsBetween(0, sys.maxsize),这会让每一行的窗口范围是从当前行到分区末尾,last取的是这个子窗口的最后值(也就是整个分区的最后值),然后直接赋值给新列,完全覆盖了原列的非空值,导致所有行都变成了分区最后值。
内容的提问来源于stack exchange,提问作者Abhra
相关产品推荐
相关产品推荐

