如何使用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
相关产品推荐
相关产品推荐

