【问题标题】:Setup a flag I or U for the Databricks delta merge为 Databricks 增量合并设置标志 I 或 U
【发布时间】:2022-01-17 01:30:40
【问题描述】:

当我执行增量合并逻辑时,有没有办法可以设置标志列(I-inserted,U-updated)。我很想知道在每日增量合并逻辑中插入了多少记录以及更新了多少记录。

我的示例数据框:

df_latest = spark.createDataFrame(
[
  ('Java', "20000"),  # create your data here, be consistent in the types.
  ('Scala', '90000'),
  ('Python', '100000')
],
["language", "users_count"]  # add your column names here
)

当我执行下面的 delta 合并逻辑时,我需要多一个名为 flag 的列(I 或 U),在 delta 表的 version02 上描述插入了多少行以及如何更新行。

test_delta.alias("h")\
  .merge(df_latest.alias("df"), "h.language = df.language")\
  .whenMatchedUpdateAll()\
  .whenNotMatchedInsertAll()\
  .execute()

任何帮助将不胜感激,我自己无法解决这个问题..!!

【问题讨论】:

    标签: apache-spark-sql delta-lake


    【解决方案1】:

    我正在寻找一种类似的技术,偶然发现了你的问题。幸运的是,我能够使用下面描述的方法解决它:

    1-) 创建一个增量表并定义它的架构。

    from delta import *
    
    DeltaTable.createIfNotExists(spark) \
        .addColumn("language",StringType())\
        .addColumn("users_count",IntegerType())\
        .addColumn("Flag",StringType())\
        .property("description", "Testing Flag Logic") \
        .location("/mnt/output/TestingFlagLogic") \
        .execute()
    

    可以在下面的链接中看到表格的快照 Snapshot of the DataFrame

    2-) 使用以下命令插入表格。

    from delta.tables import *
    
    deltaTable = DeltaTable.forPath(spark,"/mnt/output/TestingFlagLogic")
    
    deltaTable.alias("Destination")\
    .merge(
        df.alias("Updates"),
        "Destination.language = Updates.language")\
      .whenMatchedUpdate(set =
        {
          "language": "Updates.language",
          "users_count": "Updates.users_count",
          "Flag": F.lit("U")
        }) \
      .whenNotMatchedInsert( values =
        {
          "language": "Updates.language",
          "users_count": "Updates.users_count",
          "Flag": F.lit("I")
        }) \
      .execute()
    

    3-) 在第一次插入之后,你会得到下面的 DataFrame,它会有一个填充了值“I”的标志列。 Delta Table After First Insertion

    4-) 你用你想要更新的值定义了一个新的 DataFrame。这里我将“Python”和“C++”语言的用户数加倍。

    df_updated = spark.createDataFrame(
    [
      ('Python', '200000'),
      ('C++', '300000'),
    ],
    ["language", "users_count"]  # add your column names here
    )
    

    Snapshot of the Dataframe which has the values to update

    5-) 现在使用步骤 2 中描述的相同逻辑插入。只需将 df 更改为 df_updated。

    from delta.tables import *
    
    deltaTable = DeltaTable.forPath(spark,"/mnt/output/TestingFlagLogic")
    deltaTable.alias("Destination")\
    .merge(
        df_updated.alias("Updates"),
        "Destination.language = Updates.language")\
      .whenMatchedUpdate(set =
        {
          "language": "Updates.language",
          "users_count": "Updates.users_count",
          "Flag": F.lit("U")
        }) \
      .whenNotMatchedInsert( values =
        {
          "language": "Updates.language",
          "users_count": "Updates.users_count",
          "Flag": F.lit("I")
        }) \
      .execute()
    

    6-) 恭喜,您已成功实现上述功能。 现在查询您的增量并显示它以进行视觉验证。

    df_new = spark.read.format("delta").load("/mnt/output/TestingFlagLogic")
    display(df_new)
    

    更新后增量表的快照可以在下面的链接中看到。 Snapshot of the Updated Table

    您可以看到“Python”和“C++”具有更新的用户计数值以及“标志”值作为“U”,因为它是一个 Upsert(Update+Insert) 操作。

    【讨论】:

      【解决方案2】:

      如果您只需要指标,那么您可以从表的历史记录中检索该信息(通过DESCRIBE HISTORY SQL 命令或通过history function)。他们都返回了包含operation metrics 的数据帧(operationMetrics 列是一个映射),对于 MERGE 操作,有一些指标描述了插入/更新/删除的行数:numTargetRowsInsertednumTargetRowsUpdated , numTargetRowsDeleted.

      类似这样的内容(仅适用于最新版本 - 如果您需要全部,则只需从历史调用中删除 1):

      from delta.tables import *
      
      deltaTable = DeltaTable.forName(spark, "mytable")
      df = deltaTable.history(1)
      f.select(df["operationMetrics"]["numTargetRowsInserted"],
               df["operationMetrics"]["numTargetRowsUpdated"],
               df["operationMetrics"]["numTargetRowsDeleted"])
      

      【讨论】:

      • 感谢@Alex Ott,但由于某种原因它不起作用,我在 operationMetrics[numTargetRowsInserted]、operationMetrics[numTargetRowsUpdated]、operationMetrics[numTargetRowsDeleted] 上得到“null”
      • 你从df = deltaTable.history(1) 得到什么 - df.show(truncate=False) 返回什么?
      • df = deltaTable.history(1) -- 获取 version=05 行。但 operationMetrics 属性值为 NULL。
      • 您好@AlexOtt,我看不到更新、插入和删除数量的原因是因为我使用的是 Databricks 运行时版本 6.4。我需要在运行时版本> 6.4。我已经测试了在 7.2 上运行的集群,在那里我可以看到更新、删除和插入的数量,非常感谢您的帮助。
      猜你喜欢
      • 2016-04-18
      • 2021-09-18
      • 1970-01-01
      • 2023-01-31
      • 2019-02-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-18
      相关资源
      最近更新 更多