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

Delta Lake合并每日数据时如何让新版本不保留源端已删除记录

Delta Lake每日时序加载逻辑调整方案

问题说明

现有每日时序数据加载到Delta表dst1的逻辑,仅能实现新增记录插入、已有记录更新的能力,无法处理当日源数据中不存在的历史记录的自动删除需求:

  • 第1天源数据包含id=1、2、4、6四条记录,写入后生成Delta表版本0
  • 第2天源数据仅包含id=1的更新后记录,执行原有合并逻辑后生成的版本1仍然保留了id=2、4、6三条冗余记录
  • 预期版本1仅保留当日源数据存在的id=1的记录

原因分析

原有MERGE合并语句仅配置了两个处理分支:

  • 目标表与源数据id匹配时更新全字段
  • 目标表无对应id时插入源数据全字段
    对「目标表存在但源数据中无对应id」的场景没有配置处理规则,因此冗余历史记录会持续保留。

解决方案

利用Delta Lake MERGE语法的WHEN NOT MATCHED BY SOURCE分支,新增源端无匹配记录时的删除逻辑即可。

调整后完整代码

# 建表逻辑(优化增加IF NOT EXISTS判断避免重复建表报错)
spark.sql(f"CREATE TABLE IF NOT EXISTS {dtable} USING DELTA LOCATION '{dmount1}'")

# 新增删除分支的合并逻辑
spark.sql(f"""
MERGE INTO {dtable} d 
USING df ds 
ON {jkey} 
WHEN MATCHED THEN UPDATE SET * 
WHEN NOT MATCHED THEN INSERT *
WHEN NOT MATCHED BY SOURCE THEN DELETE
""")

扩展说明

如果业务需要保留历史全量数据仅做逻辑删除,可以新增is_deleted标识字段,将WHEN NOT MATCHED BY SOURCE分支的动作修改为UPDATE SET is_deleted = 1即可,不需要物理删除记录。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 20:15:02