PySpark补全DataFrame缺失日期并填充相邻值最小值
补全PySpark DataFrame缺失日期并按相邻有效Qty最小值填充
以下是实现需求的完整步骤和代码,基于PySpark 3.x版本:
1. 构造测试数据
先创建一个模拟的DataFrame,包含ID、Date和Qty列,存在日期缺失:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("FillMissingDateQty").getOrCreate() # 测试数据 data = [ ("A", "2023-10-01", 10), ("A", "2023-10-03", 20), ("A", "2023-10-06", 15), ("B", "2023-10-02", 5), ("B", "2023-10-05", 8) ] df = spark.createDataFrame(data, ["ID", "Date", "Qty"]) df = df.withColumn("Date", F.to_date("Date"))
2. 生成每个ID的完整日期序列
按ID分组,计算每个ID的日期范围,然后生成该范围内的所有连续日期:
# 计算每个ID的最小、最大日期 id_date_range = df.groupBy("ID").agg( F.min("Date").alias("min_date"), F.max("Date").alias("max_date") ) # 生成连续日期序列并展开 full_dates = id_date_range.withColumn( "Date", F.explode(F.sequence(F.col("min_date"), F.col("max_date"), F.expr("interval 1 day"))) ).select("ID", "Date")
3. 左连接原表得到包含缺失日期的DataFrame
将完整日期表与原表左连接,保留所有日期,缺失的Qty会显示为null:
df_full = full_dates.join(df, on=["ID", "Date"], how="left")
4. 用相邻有效Qty的最小值填充缺失值
通过窗口函数分别获取每个缺失日期的前一个有效Qty和后一个有效Qty,再取两者的最小值填充:
# 定义窗口:按ID分组,按日期排序 window_spec = Window.partitionBy("ID").orderBy("Date") # 向前取最近的有效Qty prev_valid_qty = F.last("Qty", ignorenulls=True).over(window_spec.rowsBetween(Window.unboundedPreceding, Window.currentRow)) # 向后取最近的有效Qty next_valid_qty = F.first("Qty", ignorenulls=True).over(window_spec.rowsBetween(Window.currentRow, Window.unboundedFollowing)) # 填充缺失值:原Qty不为空则保留,否则取前后有效Qty的最小值 df_final = df_full.withColumn( "Qty", F.when(F.col("Qty").isNotNull(), F.col("Qty")) .otherwise(F.least(prev_valid_qty, next_valid_qty)) )
查看结果
执行以下代码查看最终结果:
df_final.orderBy("ID", "Date").show()
输出示例:
+---+----------+---+ | ID| Date|Qty| +---+----------+---+ | A|2023-10-01| 10| | A|2023-10-02| 10| | A|2023-10-03| 20| | A|2023-10-04| 15| | A|2023-10-05| 15| | A|2023-10-06| 15| | B|2023-10-02| 5| | B|2023-10-03| 5| | B|2023-10-04| 5| | B|2023-10-05| 8| +---+----------+---+
关键逻辑说明
sequence函数用于生成指定日期区间内的连续日期,确保每个ID的日期无缺失;last(..., ignorenulls=True)和first(..., ignorenulls=True)分别实现向前、向后获取最近的有效Qty值;least函数用于取两个相邻有效Qty的最小值,作为缺失日期的填充值。
内容的提问来源于stack exchange,提问作者NNM
相关产品推荐
相关产品推荐

