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

修改代码以支持Databricks中Delta Lake表的删除操作

实现Delta表全量同步(含删除目标表无匹配记录)

在Databricks SQL及Databricks Runtime 12.1+版本中,可通过WHEN NOT MATCHED BY SOURCE子句处理目标表中无对应源表记录的数据,Databricks建议添加可选条件子句避免全表重写。

现有代码仅能同步新增和更新行到目标表,以下是修改后的代码,实现源表内容覆盖目标表并删除目标表中无匹配记录的需求:

try:
    # 执行Merge操作同步数据
    if allowDuplicates == "true":
        (deltadf.alias("t")
         .merge(
             partdf.alias("s"),
             f"s.primary_key_hash = t.primary_key_hash")
         .whenNotMatchedInsertAll()
         # 删除目标表中无对应源表记录的数据
         .whenNotMatchedBySourceDelete()
         .execute()
        )
    else:
        (deltadf.alias("t")
         .merge(
             partdf.alias("s"),
             "s.primary_key_hash = t.primary_key_hash")
         .whenMatchedUpdateAll("s.change_key_hash <> t.change_key_hash")
         .whenNotMatchedInsertAll()
         # 删除目标表中无对应源表记录的数据
         .whenNotMatchedBySourceDelete()
         .execute()
        )
  
    action = f"Merged into Existing Delta for table {entityName} at path: {saveloc}"

    spark.sql(f"OPTIMIZE {stageName}{regName}")  
except Exception as e:
    print(e)

关键修改说明

  • 在两个分支的Merge逻辑中均添加了.whenNotMatchedBySourceDelete()子句,该子句会删除目标表(t)中所有在源表(s)找不到匹配primary_key_hash的记录,实现源表对目标表的全量覆盖。
  • 若需避免全表扫描/重写,可给whenNotMatchedBySourceDelete()添加过滤条件,比如只删除特定时间范围的无匹配记录:
    .whenNotMatchedBySourceDelete("t.last_updated < current_date() - 7")
    
  • 确保当前Databricks Runtime版本≥12.1,否则该子句会报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 04:17:27