PySpark如何根据pad_change列区间规则填充Pad_ID列空值
PySpark实现Pad_ID向下填充方案
这个需求属于典型的前向非空值填充场景,直接用PySpark的窗口函数+last算子即可实现,核心是利用last函数的ignorenulls参数跳过空值,取最近一次出现的非空Pad_ID值。
实现步骤
- 第一步:导入依赖
from pyspark.sql import Window import pyspark.sql.functions as F
- 第二步:定义窗口规范
按TOOL_ID分区(多设备数据互不干扰),按DATE升序排序保证时间顺序,窗口范围从分区开头到当前行:
fill_window = Window.partitionBy("TOOL_ID") \ .orderBy("DATE") \ .rowsBetween(Window.unboundedPreceding, 0)
- 第三步:执行填充逻辑
df_result = df.withColumn("Pad_ID", F.last(F.col("Pad_ID"), ignorenulls=True).over(fill_window))
注意事项
- 如果同日期下存在
pad_change=1和pad_change=0的行,排序字段可以额外加一个行号标识,保证pad_change=1的行排在同日期前面,避免填充错误。 - 该方案为纯窗口算子操作,仅需一次分区shuffle,性能优异,适合大数据量场景。
内容的提问来源于stack exchange,提问作者thentangler
相关产品推荐
相关产品推荐

