【问题标题】:update the delta table in databricks with adding value to existing columns通过向现有列添加值来更新数据块中的增量表
【发布时间】:2021-08-11 07:51:46
【问题描述】:

我有一段 scala 代码,它将在 3 个不同阶段对 id_noidentifier 的信号进行计数。 代码的输出如下所示。

+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|id_no|identifier|signal01_total|signal01_without_NaN|signal01_total_valid|signal02_total|signal02_without_NaN|signal02_total_valid|signal03_total|signal03_without_NaN|signal03_total_valid|load_timestamp              |
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|050  |ident01   |25            |23                  |20                  |45            |43                  |40                  |66            |60                  |55                  |2021-08-10T16:58:30.054+0000|
|051  |ident01   |78            |70                  |68                  |15            |14                  |14                  |10            |10                  |9                   |2021-08-10T16:58:30.054+0000|
|052  |ident01   |88            |88                  |86                  |75            |73                  |70                  |16            |13                  |13                  |2021-08-10T16:58:30.054+0000|
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+

会有超过 100 个信号,因此列数将超过 300 个。

这个数据帧被写入增量表位置,如下所示。

statisticsDf.write.format("delta").option("mergeSchema", "true").mode("append").partitionBy("id_no").save(statsDestFolderPath)

对于下周的数据,我再次执行此代码并获取如下所示的数据。

+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|id_no|identifier|signal01_total|signal01_without_NaN|signal01_total_valid|signal02_total|signal02_without_NaN|signal02_total_valid|signal03_total|signal03_without_NaN|signal03_total_valid|load_timestamp              |
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|050  |ident01   |10            |8                   |7                   |15            |15                  |14                  |38            |38                  |37                  |2021-08-10T16:58:30.054+0000|
|051  |ident01   |10            |10                  |9                   |16            |15                  |15                  |30            |30                  |30                  |2021-08-10T16:58:30.054+0000|
|052  |ident01   |26            |24                  |24                  |24            |23                  |23                  |40            |38                  |36                  |2021-08-10T16:58:30.054+0000|
|053  |ident01   |25            |24                  |23                  |20            |19                  |19                  |25            |25                  |24                  |2021-08-10T16:58:30.054+0000|
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+

但我期望的输出是如果 id_noidentifiersignal name 已经存在于表中,那么它应该添加以现有数据计数,如果 id_noidentifiersignal name 是新的,则应将其添加到最终表中。

我现在收到的输出如下所示,每次运行都会附加数据。

+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|id_no|identifier|signal01_total|signal01_without_NaN|signal01_total_valid|signal02_total|signal02_without_NaN|signal02_total_valid|signal03_total|signal03_without_NaN|signal03_total_valid|load_timestamp              |
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|050  |ident01   |25            |23                  |20                  |45            |43                  |40                  |66            |60                  |55                  |2021-08-10T16:58:30.054+0000|
|051  |ident01   |78            |70                  |68                  |15            |14                  |14                  |10            |10                  |9                   |2021-08-10T16:58:30.054+0000|
|052  |ident01   |88            |88                  |86                  |75            |73                  |70                  |16            |13                  |13                  |2021-08-10T16:58:30.054+0000|
|050  |ident01   |10            |8                   |7                   |15            |15                  |14                  |38            |38                  |37                  |2021-08-10T16:58:30.054+0000|
|051  |ident01   |10            |10                  |9                   |16            |15                  |15                  |30            |30                  |30                  |2021-08-10T16:58:30.054+0000|
|052  |ident01   |26            |24                  |24                  |24            |23                  |23                  |40            |38                  |36                  |2021-08-10T16:58:30.054+0000|
|053  |ident01   |25            |24                  |23                  |20            |19                  |19                  |25            |25                  |24                  |2021-08-10T16:58:30.054+0000|
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+

但我希望得到如下所示的输出。

+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|id_no|identifier|signal01_total|signal01_without_NaN|signal01_total_valid|signal02_total|signal02_without_NaN|signal02_total_valid|signal03_total|signal03_without_NaN|signal03_total_valid|load_timestamp              |
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|050  |ident01   |35            |31                  |27                  |60            |58                  |54                  |38            |38                  |37                  |2021-08-10T16:58:30.054+0000|
|051  |ident01   |88            |80                  |77                  |31            |29                  |19                  |30            |30                  |30                  |2021-08-10T16:58:30.054+0000|
|052  |ident01   |114           |102                 |110                 |99            |96                  |93                  |40            |38                  |36                  |2021-08-10T16:58:30.054+0000|
|053  |ident01   |25            |24                  |23                  |20            |19                  |19                  |25            |25                  |24                  |2021-08-10T16:58:30.054+0000|
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+

使用 upsert 命令得到提示,如下所示。

val updatesDF = ...  // define the updates DataFrame[id_no, identifier, sig01_total, sig01_NaN, sig01_final, sig02_total,.......]

DeltaTable.forPath(spark, "/data/events/")
  .as("events")
  .merge(
    updatesDF.as("updates"),
    "events.id_no = updates.id_no" &&
    "events.identifier = updates.identifier")
  .whenMatched
  .updateExpr(
    Map("sig01_total" -> "updates.sig01_total"
                      -> 
                      ->........))
  .whenNotMatched
  .insertExpr(
    Map(
      "id_no" -> "updates.id_no",
      "identifier" -> "updates.identifier",
      "sig01_total" -> "updates.sig01_total"
                    ->
                    ->
                     .....))
  .execute()

但在我的情况下,每次的列数可能会有所不同,如果将新信号添加到 id 中,那么我们必须添加相同的。如果现有 id 的信号之一不适用于当前周流程,则仅该信号值应保持不变,其余应更新。

是否有任何选项可以使用增量表合并或更新上述代码或任何其他方式来实现此要求? 任何线索表示赞赏!

【问题讨论】:

    标签: scala apache-spark azure-databricks delta


    【解决方案1】:

    问题中提到的用例,需要一个 upsert 操作。

    您可以使用 Databricks 文档进行 upsert 操作,您可以在其中编写逻辑来执行 upsert 操作。

    您可以根据表达式控制何时插入和何时更新。

    参考链接 https://docs.databricks.com/delta/delta-update.html#upsert-into-a-table-using-merge

    【讨论】:

    • 是的,这应该可以,对我来说,除了前 2 列之外,其余列都不是静态的,这些列号因每次运行而异。那么我们应该如何在那里实现这个逻辑呢?
    • 您需要通过查看up​​datesDF的架构和现有的增量表并根据列生成映射来动态准备。
    猜你喜欢
    • 2021-09-16
    • 2020-10-07
    • 1970-01-01
    • 2014-06-05
    • 1970-01-01
    • 1970-01-01
    • 2022-12-13
    • 1970-01-01
    • 2023-02-23
    相关资源
    最近更新 更多