如何用Databricks DLT CDC实现按日取最新ID的SCD处理?
解决方案:按日保留每个Id的最新记录
方案一:使用SCD Type 1(复合主键,简洁高效)
你的需求核心是每个Id每天保留一条最新记录,其实无需直接切换到SCD Type 2,通过给SCD Type 1设置复合主键就能实现:
- 在源视图中新增日期字段,提取
PublishDateTime的日期部分; - 将主键设为
["Id", "PublishDate"],同一Id、同一日期的最新记录会覆盖旧条目,不同日期的记录则作为新行保留。
修改后的代码:
from pyspark.sql.functions import to_date, col @dlt.view def users(): # 新增日期字段,拆分出日期维度 return spark.readStream.table("source_table") \ .withColumn("PublishDate", to_date(col("PublishDateTime"))) dlt.create_streaming_table("target_table") dlt.apply_changes( target = "target_table", source = "users", keys = ["Id", "PublishDate"], # 复合主键:Id+日期 sequence_by = col("PublishDateTime"), # 按时间戳判断最新记录 stored_as_scd_type = 1 )
方案二:使用SCD Type 2(需追踪历史变更轨迹)
如果需要记录每个Id的全量版本变更(比如后续同一日期有更新时,保留旧记录的失效时间),可以使用SCD Type 2。DLT会自动为目标表生成start_date和end_date字段,标记每条记录的生效/失效时间段,最新版本的end_date为NULL。
实现代码
from pyspark.sql.functions import to_date, col @dlt.view def users(): return spark.readStream.table("source_table") \ .withColumn("PublishDate", to_date(col("PublishDateTime"))) dlt.create_streaming_table("target_table") dlt.apply_changes( target = "target_table", source = "users", keys = ["Id"], sequence_by = col("PublishDateTime"), stored_as_scd_type = 2, # 可选:排除无需追踪历史的字段,降低存储开销 track_history_except_columns = ["PublishDate"] )
查询每日最新记录
需要时可通过SQL筛选每个Id每天的最新记录:
WITH ranked_records AS ( SELECT *, ROW_NUMBER() OVER (PARTITION BY Id, PublishDate ORDER BY PublishDateTime DESC) AS rn FROM target_table ) SELECT * FROM ranked_records WHERE rn = 1
方案选择建议
- 若仅需保留每日最新记录、无审计或历史追踪需求,优先选方案一,实现简单且性能更优;
- 若需要追踪记录的全量版本变化,再选择方案二。
内容的提问来源于stack exchange,提问作者Itachi07
相关产品推荐
相关产品推荐

