【发布时间】:2021-08-11 07:51:46
【问题描述】:
我有一段 scala 代码,它将在 3 个不同阶段对 id_no 和 identifier 的信号进行计数。 代码的输出如下所示。
+-----+----------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+--------------+--------------------+--------------------+----------------------------+
|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_no 、identifier 和 signal name 已经存在于表中,那么它应该添加以现有数据计数,如果 id_no、identifier 和 signal 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