修改代码以支持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
相关产品推荐
相关产品推荐

