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

如何使用Delta Live Table实现SCD Type 1?求实时实践示例

Delta Live Table 实现SCD Type 1的实时示例

1. 定义实时数据源表

首先创建模拟的实时流入数据源(实际场景可替换为Kafka、云存储文件流等):

import dlt
from pyspark.sql.functions import *

@dlt.table(
  name="customer_source_stream",
  comment="实时客户数据流入的源表"
)
def customer_source_stream():
    return (
        spark.readStream.format("rate")
        .option("rowsPerSecond", 1)
        .load()
        .withColumn("customer_id", (rand() * 10).cast("int"))
        .withColumn("customer_name", concat(lit("Customer_"), col("customer_id"), lit("_"), col("value")))
        .withColumn("email", concat(lit("customer_"), col("customer_id"), lit("@example.com")))
        .withColumn("load_timestamp", current_timestamp())
        .select("customer_id", "customer_name", "email", "load_timestamp")
    )

2. 实现SCD Type 1目标表

使用DLT的@dlt.apply_changes注解直接实现SCD Type 1,主键匹配时自动覆盖更新非主键字段:

@dlt.table(
  name="customer_scd1",
  comment="基于SCD Type 1维护的客户维度表"
)
@dlt.apply_changes(
  target="customer_scd1",
  source="customer_source_stream",
  keys=["customer_id"],
  sequence_by="load_timestamp",
  apply_as_scd_type=1
)
def customer_scd1():
    return dlt.read("customer_source_stream")

关键参数说明

  • keys=["customer_id"]:指定匹配记录的主键字段,确保同一客户的新数据能定位到对应旧记录
  • sequence_by="load_timestamp":指定排序字段,保证最新时间戳的记录优先覆盖旧数据
  • apply_as_scd_type=1:明确采用SCD Type 1策略,主键匹配时直接覆盖所有非主键字段

验证逻辑

运行DLT流水线后,当customer_source_stream出现相同customer_id的新记录时,customer_scd1中的对应记录会被新数据直接覆盖,不保留历史版本。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 20:45:42