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

如何在Databricks中针对空目标表实现SCD Type 2(附示例)

在Databricks中实现SCD Type 2(含初始空表加载与历史数据处理)

SCD Type 2的核心是保留数据全量历史版本,通过有效起始日期、有效结束日期和当前有效标记区分不同版本。以下结合你提供的样本数据,给出从空表初始加载到后续用MERGE操作维护的完整实现,基于Delta Lake(Databricks原生支持的ACID存储格式)。


1. 准备样本源数据

假设你的样本数据顺序代表时间先后(先出现的为旧版本,后出现的为新版本),先创建源数据DataFrame:

# PySpark 代码
source_data = [
    (1, "A"),
    (2, "B"),
    (1, "A2"),
    (3, "C"),
    (3, "C2")
]
source_df = spark.createDataFrame(source_data, ["ID", "Name"])

2. 创建SCD Type 2目标表

先定义包含SCD元数据字段的目标表,初始为空:

CREATE TABLE IF NOT EXISTS scd_type2_target (
    ID INT,
    Name STRING,
    start_date DATE, -- 版本生效日期
    end_date DATE,   -- 版本失效日期(未失效用9999-12-31标记)
    is_current BOOLEAN -- 是否为当前有效版本
)
USING DELTA;

3. 初始加载历史数据到空表

因为目标表为空,需要先处理源数据中的所有历史版本,生成每个版本的有效时间范围:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, current_date, lit, when, date_sub, lead

# 按业务主键ID分区,按数据出现顺序排序(模拟时间先后)
window_rank = Window.partitionBy("ID").orderBy(lit(1))
ranked_df = source_df.withColumn("row_num", row_number().over(window_rank))

# 为每个版本初始化生效/失效日期和当前标记
scd_temp_df = ranked_df.withColumn(
    "start_date", current_date()  # 若有业务时间字段,替换为实际生效时间
).withColumn(
    "end_date", lit("9999-12-31")
).withColumn(
    "is_current", lit(True)
)

# 调整旧版本的失效日期为下一个版本的生效前一天,并更新当前标记
window_next = Window.partitionBy("ID").orderBy("row_num")
scd_final_df = scd_temp_df.withColumn(
    "next_start", lead("start_date").over(window_next)
).withColumn(
    "end_date", when(col("next_start").isNotNull(), date_sub(col("next_start"), 1)).otherwise(col("end_date"))
).withColumn(
    "is_current", when(col("next_start").isNotNull(), lit(False)).otherwise(col("is_current"))
).drop("row_num", "next_start")

# 写入空目标表
scd_final_df.write.mode("append").saveAsTable("scd_type2_target")

4. 用MERGE操作处理后续追加数据

当目标表已有数据后,后续新增/变更数据通过MERGE操作自动维护SCD Type2:

示例追加数据

new_source_data = [
    (1, "A3"),  # ID=1的新版本
    (4, "D")    # 全新ID的记录
]
new_source_df = spark.createDataFrame(new_source_data, ["ID", "Name"])

执行MERGE操作

MERGE INTO scd_type2_target t
USING (
    SELECT 
        ID, 
        Name,
        current_date() AS start_date,
        CAST('9999-12-31' AS DATE) AS end_date,
        TRUE AS is_current
    FROM new_source_df
) s
ON t.ID = s.ID AND t.is_current = TRUE -- 匹配当前有效的同ID记录
WHEN MATCHED AND t.Name != s.Name THEN -- 字段发生变更
    UPDATE SET 
        t.end_date = DATE_SUB(s.start_date, 1),
        t.is_current = FALSE
WHEN NOT MATCHED THEN -- 无匹配记录,插入新版本
    INSERT (ID, Name, start_date, end_date, is_current)
    VALUES (s.ID, s.Name, s.start_date, s.end_date, s.is_current)

验证结果

查询目标表可看到所有历史版本:

SELECT * FROM scd_type2_target ORDER BY ID, start_date;

预期输出(假设当前日期为2024-05-20):

IDNamestart_dateend_dateis_current
1A2024-05-202024-05-19false
1A22024-05-202024-05-19false
1A32024-05-209999-12-31true
2B2024-05-209999-12-31true
3C2024-05-202024-05-19false
3C22024-05-209999-12-31true
4D2024-05-209999-12-31true

提示:如果源数据包含业务时间字段(如update_time),建议用该字段替换current_date(),能更准确反映数据实际生效时间。

内容的提问来源于stack exchange,提问作者CloudEngineer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 13:14:52