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

基于Delta Lake实现SCD Type 2的技术方案求助及疑问解答

在Databricks中用PySpark实现Delta表的SCD Type2

疑问解答

1. 为何需要单独的代理键与合并键?

  • 合并键:是识别源表和目标表中同一业务实体的字段组合(比如这里的state, Code, Name),核心作用是判断业务实体属于新增、待更新还是已从源表移除。但合并键只能定位到业务实体,无法区分同一实体的不同版本记录(比如同一state+Code+Name下,value变更后的新记录和旧记录)。
  • 代理键:是为每条记录生成的唯一标识符(与业务无关),用来区分同一业务实体的不同历史版本。比如SCD Type2中,同一业务实体会有多条记录(旧记录标记为非活跃,新记录为活跃),代理键能让你精准定位到某一条历史记录,避免业务键变更导致的关联错误,同时方便后续的报表、分析或数据关联操作。

2. 基于[state,code,name,value]生成代理键并端到端实现

我们可以用sha2哈希函数基于state, code, name, value生成固定长度的字符串作为代理键(避免分布式环境下自增ID的冲突问题),下面是完整的实现步骤:


端到端实现代码

1. 模拟源表数据

from pyspark.sql.functions import current_timestamp, sha2, concat_ws, lit

# 模拟源表数据(包含新增、更新、缺失的记录)
source_data = [
    ("CA", "CA001", "California", 100, current_timestamp()),
    ("NY", "NY001", "New York", 200, current_timestamp()),
    ("BY", "BY001", "Bayern", 150, current_timestamp()),  # value从120更新为150
    ("TX", "TX001", "Texas", 300, current_timestamp())   # 新增记录
]

source_df = spark.createDataFrame(
    source_data,
    schema=["state", "Code", "Name", "value", "insert_datetime"]
)

# 生成代理键:基于state,code,name,value的哈希值
source_df = source_df.withColumn(
    "surrogate_key",
    sha2(concat_ws("|", "state", "Code", "Name", "value"), 256)
)

2. 初始化目标Delta表(如果不存在)

# 定义目标表schema
target_schema = """
    state STRING,
    Code STRING,
    Name STRING,
    value INT,
    insert_datetime TIMESTAMP,
    is_current BOOLEAN,
    ExpiryDate TIMESTAMP,
    surrogate_key STRING
"""

# 初始化空的Delta表(如果不存在)
if not spark.catalog.tableExists("default.silver_scd2_table"):
    spark.createDataFrame([], schema=target_schema).write.format("delta").mode("overwrite").saveAsTable("default.silver_scd2_table")

3. 执行SCD Type2合并逻辑

from delta.tables import DeltaTable

# 加载目标Delta表
target_table = DeltaTable.forName(spark, "default.silver_scd2_table")

# 定义合并条件:匹配同一业务实体(state+Code+Name),且目标表记录为当前活跃
merge_condition = """
    target.state = source.state AND
    target.Code = source.Code AND
    target.Name = source.Name AND
    target.is_current = true
"""

# 第一步:处理新增和更新记录
target_table.alias("target").merge(
    source_df.alias("source"),
    merge_condition
).whenMatchedUpdate(
    # 当匹配到且value不同时,标记旧记录为非活跃,设置过期时间
    condition="target.value != source.value",
    set={
        "is_current": lit(False),
        "ExpiryDate": current_timestamp()
    }
).whenNotMatchedInsert(
    # 插入新记录:标记为活跃,过期时间设为null
    values={
        "state": "source.state",
        "Code": "source.Code",
        "Name": "source.Name",
        "value": "source.value",
        "insert_datetime": "source.insert_datetime",
        "is_current": lit(True),
        "ExpiryDate": lit(None),
        "surrogate_key": "source.surrogate_key"
    }
).execute()

# 第二步:处理源表缺失的记录(软删除)
# 先获取源表的业务实体集合
source_business_keys = source_df.select("state", "Code", "Name").distinct()

# 更新目标表中不在源集合的活跃记录,标记为非活跃
target_table.alias("target").merge(
    source_business_keys.alias("source"),
    """
        target.state = source.state AND
        target.Code = source.Code AND
        target.Name = source.Name
    """
).whenNotMatchedUpdate(
    condition="target.is_current = true",
    set={
        "is_current": lit(False),
        "ExpiryDate": current_timestamp()
    }
).execute()

4. 验证结果

# 查询目标表,查看SCD Type2效果
spark.sql("SELECT * FROM default.silver_scd2_table ORDER BY state, is_current DESC").show(truncate=False)

关键说明

  • 代理键生成:用sha2(concat_ws("|", ...), 256)确保同一state+Code+Name+value组合生成唯一哈希值,适配Databricks分布式环境。
  • 删除处理:采用软删除方式保留历史记录,符合SCD Type2的特性;若业务要求硬删除,可替换为whenNotMatchedDelete(不推荐,会丢失历史数据)。
  • 合并逻辑:分两步执行,先处理新增和更新,再处理源表缺失的活跃记录,确保逻辑清晰。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 15:38:21