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

PySpark按规则填充Interest_rate空值:缺失类型后Fixed值异常

问题:Spark DataFrame按规则填充空值异常

原始DataFrame

IdDateInt_typeInterest_rate
A03/22/2023Floating0.044
A03/22/2023Floating0.045
A03/22/2023Floating0.046
A03/24/2023Floating0.046
A03/24/2023FixedNull
A03/24/2023FixedNull
A03/24/2023MissingNull
A03/24/2023MissingNull
A03/24/2023FixedNull
A03/24/2023FixedNull
A03/24/2023FixedNull
A03/24/2023FixedNull
A03/24/2023FixedNull

需求规则

  • 当Interest_rate为空,且同Id、Date分组内的前一行有可用值,且当前行Int_type为FIXED时,向前填充该利率
  • 当Int_type为Missing时,Interest_rate默认保持Null

尝试的代码

wind1=Window.partitionBy('id','date').orderBy('date')

df = df.withColumn('lag_int_rt',when((upper(col('Int_Type'))=='FIXED') & (col('Interest_rate').isNull()),lag('Interest_rate').over(wind1)))

df = df.withColumn('lag_int_rt',when((upper(col('Int_Type'))=='FIXED') & (col('lag_int_rt').isNull()) & (upper(lag('Int_Type').over(wind1)) !='MISSING') ,\
last('lag_int_rt',True).over(wind1)).otherwise(col('lag_int_rt')))

df = df.withColumn('final_interest',coalesce('Interest_rate','lag_int_rt'))

当前问题

现有代码处理后,最后5行Int_type为Fixed的行被错误填充了0.046。按规则,这些行的前一行是Missing类型,应该保持Null,但代码仅第一个Fixed行处理正确,后续连续Fixed行未按规则更新。

期望输出

IdDateInt_typeInterest_rateFinal_interest
A03/22/2023Floating0.0440.044
A03/22/2023Floating0.0450.045
A03/22/2023Floating0.0460.046
A03/24/2023Floating0.0460.046
A03/24/2023FixedNull0.046
A03/24/2023FixedNull0.046
A03/24/2023MissingNullNull
A03/24/2023MissingNullNull
A03/24/2023FixedNullNull
A03/24/2023FixedNullNull
A03/24/2023FixedNullNull
A03/24/2023FixedNullNull
A03/24/2023FixedNullNull

解决方案

原代码未通过Missing行阻断跨块填充,导致后续Fixed行错误复用了之前块的有效值。正确做法是先按Missing行将同Id+Date的数据拆分成分组块,再在块内进行填充:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 1. 按Id、Date分区,标记Missing行并生成分组块ID
window_part = Window.partitionBy("Id", "Date").orderBy(F.monotonically_increasing_id())
df = df.withColumn(
    "is_missing",
    F.when(F.upper(F.col("Int_type")) == "MISSING", 1).otherwise(0)
).withColumn(
    "block_id",
    F.sum("is_missing").over(window_part)
)

# 2. 在每个Id+Date+block_id块内,获取最近的非空Interest_rate
window_fill = Window.partitionBy("Id", "Date", "block_id").orderBy(F.monotonically_increasing_id())
df = df.withColumn(
    "last_valid_rate",
    F.last(F.when(F.col("Interest_rate").isNotNull(), F.col("Interest_rate")), ignorenulls=True).over(window_fill)
)

# 3. 生成最终的final_interest
df = df.withColumn(
    "final_interest",
    F.when(
        F.upper(F.col("Int_type")) == "MISSING",
        None
    ).when(
        F.upper(F.col("Int_type")) == "FIXED" & F.col("Interest_rate").isNull(),
        F.col("last_valid_rate")
    ).otherwise(F.col("Interest_rate"))
).drop("is_missing", "block_id", "last_valid_rate")

# 查看结果
df.show()

代码说明

  • 第一步通过is_missing标记Missing行,再用累加和生成block_id,将每个Missing行之后的数据分到新块,实现Missing行的阻断效果
  • 第二步在每个块内,用last函数仅查找当前块内最近的非空利率,不会跨Missing行取值
  • 第三步严格按照规则生成最终值:Missing行直接设为Null,Fixed空值行用块内有效值,其他行保留原值

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 12:55:02