如何在SQL/PySpark中补全日期区间内缺失的学生余额数据
补全缺失日期并沿用前值的SQL/PySpark实现
问题背景
现有一张按DATE列分区的学生资产组合余额表,表结构如下:
NUM1 STRING (Examp: 2343) NUM2 STRING (Examp: 0982) DOC STRING (Examp: 0987654321) CLASS STRING (Examp: RED / BLACK - 019 ) COD_CLASS STRING (Examp: 9087) -- 修正原表中"COD CLASS"的字段名格式 NOME_CLASS STRING (Examp: REDBCK ) DATE STRING (Examp: 09-10-2022 ) BALANCE STRING (Examp: 10,00 )
源数据示例:
2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 24-02-2023, 90,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 26-02-2023, 00,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 02-03-2023, 80,02
需求:补全2023年2月24日至2023年3月2日之间的缺失日期,缺失日期的BALANCE值沿用前一个存在日期的对应值,目标结果如下:
2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 24-02-2023, 90,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 25-02-2023, 90,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 26-02-2023, 00,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 27-02-2023, 00,02 2343, 0982, 0987654321, RED / BLACK , 9087, REDBCK, 02-03-2023, 80,02
SQL实现(Spark SQL为例)
核心思路:生成目标日期范围的完整日期序列,与学生唯一标识组合交叉关联,再通过窗口函数向前填充缺失的BALANCE值。
-- 1. 生成目标日期范围的临时表 WITH date_series AS ( SELECT date_add(to_date('2023-02-24'), pos) AS target_date FROM posexplode(array_repeat(0, datediff(to_date('2023-03-02'), to_date('2023-02-24')) + 1)) ), -- 2. 获取学生唯一标识分组 student_groups AS ( SELECT DISTINCT NUM1, NUM2, DOC, CLASS, COD_CLASS, NOME_CLASS FROM your_table ), -- 3. 交叉关联生成每个学生的所有日期 student_dates AS ( SELECT sg.*, ds.target_date FROM student_groups sg CROSS JOIN date_series ds ), -- 4. 左连接原表匹配已有数据 joined_data AS ( SELECT sd.NUM1, sd.NUM2, sd.DOC, sd.CLASS, sd.COD_CLASS, sd.NOME_CLASS, sd.target_date, t.BALANCE FROM student_dates sd LEFT JOIN your_table t ON sd.NUM1 = t.NUM1 AND sd.NUM2 = t.NUM2 AND sd.DOC = t.DOC AND to_date(t.DATE, 'dd-MM-yyyy') = sd.target_date ) -- 5. 用LAST_VALUE向前填充,输出结果 SELECT NUM1, NUM2, DOC, CLASS, COD_CLASS, NOME_CLASS, date_format(target_date, 'dd-MM-yyyy') AS DATE, LAST_VALUE(BALANCE, TRUE) OVER ( PARTITION BY NUM1, NUM2, DOC, CLASS, COD_CLASS, NOME_CLASS ORDER BY target_date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS BALANCE FROM joined_data WHERE target_date BETWEEN to_date('2023-02-24') AND to_date('2023-03-02') ORDER BY target_date;
PySpark实现
核心思路:生成日期序列DataFrame,与学生分组交叉连接,左连接原表后使用last函数向前填充缺失值。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, date_add, to_date, date_format, last from pyspark.sql.window import Window spark = SparkSession.builder.appName("FillMissingDates").getOrCreate() # 1. 读取原表并转换日期格式 df = spark.read.table("your_table") df = df.withColumn("date_dt", to_date(col("DATE"), "dd-MM-yyyy")) # 2. 生成目标日期范围的DataFrame start_date = to_date("2023-02-24") end_date = to_date("2023-03-02") date_diff = spark.sql(f"SELECT datediff({end_date}, {start_date}) AS diff").collect()[0]["diff"] date_series = spark.range(0, date_diff + 1).select(date_add(start_date, col("id")).alias("target_date")) # 3. 获取学生唯一标识分组 student_groups = df.select("NUM1", "NUM2", "DOC", "CLASS", "COD_CLASS", "NOME_CLASS").distinct() # 4. 交叉连接生成每个学生的所有日期 student_dates = student_groups.crossJoin(date_series) # 5. 左连接原表匹配已有数据 joined_df = student_dates.join( df, (student_dates.NUM1 == df.NUM1) & (student_dates.NUM2 == df.NUM2) & (student_dates.DOC == df.DOC) & (student_dates.target_date == df.date_dt), how="left" ).select( student_dates.NUM1, student_dates.NUM2, student_dates.DOC, student_dates.CLASS, student_dates.COD_CLASS, student_dates.NOME_CLASS, student_dates.target_date, df.BALANCE ) # 6. 定义窗口并向前填充BALANCE window_spec = Window.partitionBy("NUM1", "NUM2", "DOC", "CLASS", "COD_CLASS", "NOME_CLASS")\ .orderBy("target_date")\ .rowsBetween(Window.unboundedPreceding, Window.currentRow) result_df = joined_df.withColumn( "BALANCE", last("BALANCE", ignorenulls=True).over(window_spec) ).withColumn( "DATE", date_format(col("target_date"), "dd-MM-yyyy") ).select( "NUM1", "NUM2", "DOC", "CLASS", "COD_CLASS", "NOME_CLASS", "DATE", "BALANCE" ).orderBy("target_date") # 查看结果 result_df.show()
内容的提问来源于stack exchange,提问作者Felipe
相关产品推荐
相关产品推荐

