【发布时间】:2019-03-30 23:36:20
【问题描述】:
我有一个静态DataFrame,包含数百万行,如下所示。
静态DataFrame:
--------------
id|time_stamp|
--------------
|1|1540527851|
|2|1540525602|
|3|1530529187|
|4|1520529185|
|5|1510529182|
|6|1578945709|
--------------
现在,在每个批次中,正在形成一个 Streaming DataFrame,其中包含 id 和经过如下操作后更新的 time_stamp。
第一批:
--------------
id|time_stamp|
--------------
|1|1540527888|
|2|1540525999|
|3|1530529784|
--------------
现在,在每个批次中,我都想使用 Streaming Dataframe 的更新值来更新 Static DataFrame,如下所示。 怎么做?
第一批后的静态DF:
--------------
id|time_stamp|
--------------
|1|1540527888|
|2|1540525999|
|3|1530529784|
|4|1520529185|
|5|1510529182|
|6|1578945709|
--------------
我已经尝试过 except()、union() 或 'left_anti' join。但似乎结构化流不支持此类操作。
【问题讨论】:
-
如果我没记错的话,不仅结构化流,而且大多数流框架,当流和静态表连接发生时,状态表仅用作查找表,而不是插入/更新的东西.在这种情况下,Stream 将主要读取和更新数据源。因此,除了结构化流式传输之外,您可能还需要一些技巧(不确定是否有)。
-
嗨@JungtaekLim。感谢回复。这是一个特殊的情况,我必须更新静态数据帧,因为它正在与流数据帧的其他部分一起使用。
-
你好@Swarup,你找到任何方法了吗?
-
您好@Allan,如果您使用的是 Spark 版本 > 2.4.0,您只需要在数据帧 writestream 上调用 foreachBatch((batch: DataFrame, batchId: Long)。但是,我仍然不知道如何在旧版本的 Spark 上执行此操作。但是您为什么不利用新版本呢?
标签: apache-spark apache-spark-sql spark-structured-streaming