【问题标题】:How to update a Static Dataframe with Streaming Dataframe in Spark structured streaming如何在 Spark 结构化流中使用流数据帧更新静态数据帧
【发布时间】: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


【解决方案1】:

所以我通过 Spark 2.4.0 AddBatch 方法解决了这个问题,该方法将流数据帧转换为迷你批处理数据帧。但是对于

【讨论】:

  • 您能否详细解释一下 AddBatch 方法来克服这个问题?你能分享一些实现这个的代码吗?我现在遇到了同样的问题。
  • 嘿@Hong。 Kompe 已经在上面的评论中解释了这个过程。
【解决方案2】:

正如 Swarup 自己所解释的,如果您使用 Spark 2.4.x,则可以使用 forEachBatch 输出接收器。

接收器采用函数(batchDF: DataFrame, batchId: Long) => Unit,其中batchDF 是流数据帧的当前处理批次,可以用作静态数据帧。 因此,在此函数中,您可以使用每个批次的值更新其他数据帧。

请参见下面的示例: 假设您有一个名为 frameToBeUpdated 的数据框,其架构与实例变量相同,并且您希望将状态保留在那里

df
  .writeStream
  .outputMode("append")
  .foreachBatch((batch: DataFrame, batchId: Long) => {
   //batch is a static dataframe

      //take all rows from the original frames that aren't in batch and 
      //union them with the batch, then reassign to the
      //dataframe you want to keep
      frameToBeUpdated = batch.union(frameToBeUpdated.join(batch, Seq("id"), "left_anti"))
    })
    .start()

更新逻辑来自:spark: merge two dataframes, if ID duplicated in two dataframes, the row in df1 overwrites the row in df2

【讨论】:

    【解决方案3】:

    我也有类似的问题。下面是我申请更新静态数据帧的 foreachBatch。我想知道如何返回在 foreachBatch 中完成的更新的 df。

    def update_reference_df(df, static_df):
        query: StreamingQuery = df \
            .writeStream \
            .outputMode("append") \
            .format("memory") \
            .foreachBatch(lambda batch_df, batchId: update_static_df(batch_df, static_df)) \
            .start()
        return query
    
    def update_static_df(batch_df, static_df):
        df1: DataFrame = static_df.union(batch_df.join(static_df,
                                                     (batch_df.SITE == static_df.SITE)
                                                     "left_anti"))
    
        return df1
    

    【讨论】:

      猜你喜欢
      • 2021-10-23
      • 1970-01-01
      • 2018-03-14
      • 1970-01-01
      • 2018-12-16
      • 1970-01-01
      • 2018-11-08
      • 1970-01-01
      • 2021-12-03
      相关资源
      最近更新 更多