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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 13:35:16