【问题标题】:How to Prevent Duplicate Entries to enter to delta lake of Azure Storage如何防止重复条目进入 Azure Storage 的 delta Lake
【发布时间】:2021-06-10 12:15:10
【问题描述】:

我有一个以 delta 格式存储到 Adls 中的数据框,现在当我尝试将新的更新行附加到该 delta 湖时,它应该有什么方法可以删除 delta 中的旧现有记录并添加新的更新记录。

对于存储在 Delta 中的 DataFrame 架构,有一个唯一的 Column。通过它我们可以检查记录是更新的还是新的。

【问题讨论】:

    标签: azure databricks azure-databricks delta-lake azure-data-lake-gen2


    【解决方案1】:

    这是Merge command 的任务 - 您定义合并条件(您的唯一列)然后执行操作。在 SQL 中,它可能如下所示(column 是您的唯一列,@​​987654323@ 可能是您注册为临时视图的数据框):

    MERGE INTO destination
    USING updates
    ON destination.column = updates.column
    WHEN MATCHED THEN
      UPDATE SET *
    WHEN NOT MATCHED
      THEN INSERT *
    

    在 Python 中它可能如下所示:

    from delta.tables import *
    
    deltaTable = DeltaTable.forPath(spark, "/data/destination/")
    
    deltaTable.alias("dest").merge(
        updatesDF.alias("updates"),
        "dest.column = updates.column") \
      .whenMatchedUpdateAll() \
      .whenNotMatchedInsertAll() \
      .execute()
    

    【讨论】:

    • 嗨,Alex Ott,您的解决方案真的很有效,谢谢。有什么方法可以减少我们对增量表的考虑,例如将其更改为仅考虑特定分区而不是整个增量表。那真的很有帮助
    • 你只需要正确构造ON条件
    • 嗨,Alex,您能否就语法更简单地解释一下。我使用的是你提到的python代码。” deltaTable = DeltaTable.forPath(spark, “/data/destination/”)”,所以我们可以在这里提及任何内容,而不是读取整个表只读增量表的特定分区,因为天数通过数据将会增加。
    • 扩展条件"dest.column = updates.column" 包括对分区键的检查
    猜你喜欢
    • 2016-12-20
    • 1970-01-01
    • 1970-01-01
    • 2021-01-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多