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
相关产品推荐
相关产品推荐

