如何在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):
| ID | Name | start_date | end_date | is_current |
|---|---|---|---|---|
| 1 | A | 2024-05-20 | 2024-05-19 | false |
| 1 | A2 | 2024-05-20 | 2024-05-19 | false |
| 1 | A3 | 2024-05-20 | 9999-12-31 | true |
| 2 | B | 2024-05-20 | 9999-12-31 | true |
| 3 | C | 2024-05-20 | 2024-05-19 | false |
| 3 | C2 | 2024-05-20 | 9999-12-31 | true |
| 4 | D | 2024-05-20 | 9999-12-31 | true |
提示:如果源数据包含业务时间字段(如
update_time),建议用该字段替换current_date(),能更准确反映数据实际生效时间。
内容的提问来源于stack exchange,提问作者CloudEngineer
相关产品推荐
相关产品推荐

