PySpark按规则填充Interest_rate空值:缺失类型后Fixed值异常
问题:Spark DataFrame按规则填充空值异常
原始DataFrame
| Id | Date | Int_type | Interest_rate |
|---|---|---|---|
| A | 03/22/2023 | Floating | 0.044 |
| A | 03/22/2023 | Floating | 0.045 |
| A | 03/22/2023 | Floating | 0.046 |
| A | 03/24/2023 | Floating | 0.046 |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Missing | Null |
| A | 03/24/2023 | Missing | Null |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Fixed | Null |
| A | 03/24/2023 | Fixed | Null |
需求规则
- 当
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行未按规则更新。
期望输出
| Id | Date | Int_type | Interest_rate | Final_interest |
|---|---|---|---|---|
| A | 03/22/2023 | Floating | 0.044 | 0.044 |
| A | 03/22/2023 | Floating | 0.045 | 0.045 |
| A | 03/22/2023 | Floating | 0.046 | 0.046 |
| A | 03/24/2023 | Floating | 0.046 | 0.046 |
| A | 03/24/2023 | Fixed | Null | 0.046 |
| A | 03/24/2023 | Fixed | Null | 0.046 |
| A | 03/24/2023 | Missing | Null | Null |
| A | 03/24/2023 | Missing | Null | Null |
| A | 03/24/2023 | Fixed | Null | Null |
| A | 03/24/2023 | Fixed | Null | Null |
| A | 03/24/2023 | Fixed | Null | Null |
| A | 03/24/2023 | Fixed | Null | Null |
| A | 03/24/2023 | Fixed | Null | Null |
解决方案
原代码未通过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
相关产品推荐
相关产品推荐

