动态为PySpark DataFrame添加月末(EOM)列并填充空值
PySpark按ID分组填充连续月末日期(基于最后一个非空EOM值)
问题背景
我手头的PySpark DataFrame结构如下:
| id | EOM |
|---|---|
| 1 | 2023-01-31 |
| 1 | null |
| 1 | null |
| 2 | 2023-03-31 |
| 2 | null |
需求是按id分组,用每个id的最后一个非空EOM值为所有空值填充连续的月末日期,最终得到的目标DataFrame如下:
| id | EOM |
|---|---|
| 1 | 2023-01-31 |
| 1 | 2023-02-28 |
| 1 | 2023-03-31 |
| 2 | 2023-03-31 |
| 2 | 2023-04-30 |
我试了下面这段代码,但结果完全不符合预期:
df.where("EOM IS not NULL").groupBy(df['id']).agg(add_months(first(df['EOM']),1))
问题出在哪
你这段代码只筛选了非空EOM的行,按id分组后取第一个EOM加了一个月,但根本没关联回原数据填充空值,而且用了first而不是last来取最后一个非空EOM,完全没覆盖连续日期生成的逻辑。
正确实现方案
得用窗口函数来处理,核心思路是先拿到每个id的最后一个非空EOM,再计算每行相对于这个起始点的月份偏移,最后生成连续日期:
完整代码
from pyspark.sql import Window from pyspark.sql.functions import col, last, row_number, add_months, when # 1. 定义窗口:按id分组,把非空EOM排前面,方便取最后一个非空值 window_id = Window.partitionBy("id").orderBy(col("EOM").desc_nulls_last()) # 2. 给每行标记当前id的最后一个非空EOM,同时按顺序给行编号 df_with_base = df.withColumn( "last_non_null_eom", last(col("EOM"), ignorenulls=True).over(window_id) ).withColumn( "row_num", row_number().over(Window.partitionBy("id").orderBy(col("EOM").asc_nulls_last())) ) # 3. 找到每个id最后一个非空EOM所在的行号 window_last_row = Window.partitionBy("id") df_with_offset = df_with_base.withColumn( "last_non_null_row", last(when(col("EOM").isNotNull(), col("row_num")), ignorenulls=True).over(window_last_row) ) # 4. 填充空值:非空值保留,空值用最后一个非空EOM加偏移月份 result_df = df_with_offset.withColumn( "EOM", when( col("EOM").isNotNull(), col("EOM") ).otherwise( add_months(col("last_non_null_eom"), col("row_num") - col("last_non_null_row")) ) ).select("id", "EOM") # 查看结果 result_df.show()
代码说明
window_id:按id分组后,把非空EOM排在前面,这样last函数能精准取到每个id的最后一个非空EOM值。row_num:给每个id下的行按顺序编号,用来计算当前行和最后一个非空EOM行的差值,这个差值就是需要往后加的月份数。last_non_null_row:定位每个id最后一个非空EOM所在的行号,用来计算偏移量。- 最后一步的
when-otherwise:非空值直接保留,空值用最后一个非空EOM加上偏移月份,add_months会自动处理不同月份的月末天数(比如2月的28/29日)。
内容的提问来源于stack exchange,提问作者code_bug
相关产品推荐
相关产品推荐

