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

如何用Databricks DLT CDC实现按日取最新ID的SCD处理?

解决方案:按日保留每个Id的最新记录

方案一:使用SCD Type 1(复合主键,简洁高效)

你的需求核心是每个Id每天保留一条最新记录,其实无需直接切换到SCD Type 2,通过给SCD Type 1设置复合主键就能实现:

  1. 在源视图中新增日期字段,提取PublishDateTime的日期部分;
  2. 将主键设为["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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 08:08:17