【问题标题】:Restart streaming query without stopping application在不停止应用程序的情况下重新启动流式查询
【发布时间】:2023-03-29 13:22:01
【问题描述】:

我尝试使用下面的代码代替 query.awaitTermination() 在 spark 中重新启动流式查询,下面的代码将在无限循环中并寻找触发器以重新启动查询,然后执行下面的代码。基本上我试图刷新缓存的 df .

 query.processAllavaialble()
    query.stop()
   //oldDF is a cached Dataframe created from GlobalTempView which is of size 150GB.
          oldDF.unpersist()
    val inputDf: DataFrame = readFile(spec, sparkSession) //read file from S3
    or anyother source
    val recreateddf = inputDf.persist()
    //Start the query// here should i start query again by invoking readStream ?

但是当我查看 spark 文档时,它说

void processAllAvailable() ///documentation says This method is intended for testing/// Blocks until all available data in the source has been processed and committed to the sink. This method is intended for testing. Note that in the case of continually arriving data, this method may block forever. Additionally, this method is only guaranteed to block until data that has been synchronously appended data to a Source prior to invocation. (i.e. getOffset must immediately reflect the addition).


stop() Stops the execution of this query if it is running. This method blocks until the threads performing execution has stopped.

那么在不停止我的 Spark 流应用程序的情况下重新启动查询的更好方法是什么

【问题讨论】:

  • 我在答案中添加了一个参考示例

标签: apache-spark spark-streaming


【解决方案1】:

这对我有用。

下面是我在 spark 2.4.5 中针对左外连接和左连接的场景。下面的过程是推动 spark 读取最新的维度数据变化。

流程用于批量维度的流连接(始终更新)

第 1 步:-

在开始 Spark 流作业之前:- 确保维度批处理数据文件夹只有一个文件,并且该文件应该至少有一条记录(由于某种原因放置空文件不起作用)/

第 2 步:- 开始您的流媒体作业并在 kafka 流中添加流记录

第 3 步:- 用值覆盖 dim 数据(文件名应保持不变,维度文件夹应只有一个文件) 注意:- 不要使用 spark 写入此文件夹,请使用 Java 或 Scala filesystem.io 覆盖文件或 bash 删除文件并替换为具有相同名称的新数据文件。

第 4 步:- 在下一批中,spark 能够在加入 kafka 流时读取更新的维度数据...

示例代码:-

package com.databroccoli.streaming.streamjoinupdate

import org.apache.log4j.{Level, Logger}
import org.apache.spark.sql.types.{StringType, StructField, StructType, TimestampType}
import org.apache.spark.sql.{DataFrame, SparkSession}

object BroadCastStreamJoin3 {

  def main(args: Array[String]): Unit = {
    @transient lazy val logger: Logger = Logger.getLogger(getClass.getName)

    Logger.getLogger("akka").setLevel(Level.WARN)
    Logger.getLogger("org").setLevel(Level.ERROR)
    Logger.getLogger("com.amazonaws").setLevel(Level.ERROR)
    Logger.getLogger("com.amazon.ws").setLevel(Level.ERROR)
    Logger.getLogger("io.netty").setLevel(Level.ERROR)

    val spark = SparkSession
      .builder()
      .master("local")
      .getOrCreate()

    val schemaUntyped1 = StructType(
      Array(
        StructField("id", StringType),
        StructField("customrid", StringType),
        StructField("customername", StringType),
        StructField("countrycode", StringType),
        StructField("timestamp_column_fin_1", TimestampType)
      ))

    val schemaUntyped2 = StructType(
      Array(
        StructField("id", StringType),
        StructField("countrycode", StringType),
        StructField("countryname", StringType),
        StructField("timestamp_column_fin_2", TimestampType)
      ))

    val factDf1 = spark.readStream
      .schema(schemaUntyped1)
      .option("header", "true")
      .csv("src/main/resources/broadcasttest/fact")


    val dimDf3 = spark.read
      .schema(schemaUntyped2)
      .option("header", "true")
      .csv("src/main/resources/broadcasttest/dimension")
      .withColumnRenamed("id", "id_2")
      .withColumnRenamed("countrycode", "countrycode_2")

    import spark.implicits._

    factDf1
      .join(
        dimDf3,
        $"countrycode_2" <=> $"countrycode",
        "inner"
      )
      .writeStream
      .format("console")
      .outputMode("append")
      .start()
      .awaitTermination

  }
}

【讨论】:

    【解决方案2】:

    您的问题有点不清楚(第二段代码没有使用您想要保留的 df 所以我不确定您打算如何集成它们......我假设加入?

    我们遇到了类似的问题(使用 Spark 2.1),并通过创建 Sink (https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/Sink.scala) 的自定义实现来解决它,其中数据在 addBatch 中加载。由于您的设置表明您一次只处理 1 个文件并且没有水印,您可能可以将您的逻辑塞进 addBatch 方法中......虽然这有点 hacky(我相信我们的用例略有不同)。

    如果 spark 2.2 是一个选项,那么你很幸运。 Spark 2.2 添加了“运行一次”触发器,允许您将 Spark Streaming API 用于批处理作业(这实际上是您正在尝试做的事情)。如果您修改写入流以使用这个新触发器,那么无限循环可能会起作用(尽管我从未尝试过)。您最好使用外部调度程序以批处理模式运行流作业。您可以在此处阅读有关 Run Once 触发器的更多信息:https://databricks.com/blog/2017/05/22/running-streaming-jobs-day-10x-cost-savings.html

    如果您使用的是 EMR,那么 Spark 2.2 尚不可用...但我听说它将在接下来的几周内发布(祈祷)。

    您可以在这里找到一些完整的 Sink 实现示例:https://github.com/holdenk/spark-structured-streaming-ml/blob/master/src/main/scala/com/high-performance-spark-examples/structuredstreaming/CustomSink.scala

    【讨论】:

    • 琼斯:对不起,如果我的问题有点令人困惑,但我试图解决的问题陈述是在不停止应用程序的情况下刷新缓存的数据帧..当我查看 spark 邮件列表时,TDas 建议重新启动流式查询对于类似的问题。所以我想知道如何重新启动流查询
    • 我正在使用 Spark 2.1,您能否通过示例告诉我 addbatch 如何在不停止流应用程序的情况下刷新缓存的数据帧有用
    • spark 流不是真正的“流式传输”,它是微批处理。当您使用 writeStream 方法时,您将创建一个 DataStreamWriter 对象。 DataStreamWriter 的 'format' 方法用于指定要使用的输出源。实现 Sink 允许您创建这样一个新的输出源。 addBatch 方法采用 dataFrame、sparkSession 和选项,并在每个微批处理上执行 addBatch。这意味着 addBatch 中的代码会执行每个微批处理。主要缺点是所有下游转换也必须在 addBatch 方法中进行。稍后我会尝试添加示例。
    • 对不起,我有点困惑,是否应该在 addbatch 方法中添加重新启动查询..试图查看如何在 addbatch 中刷新缓存的数据帧?我了解您编写自定义接收器的建议,但它与我正在尝试解决的问题
    猜你喜欢
    • 2016-12-03
    • 1970-01-01
    • 1970-01-01
    • 2018-08-22
    • 2017-05-03
    • 1970-01-01
    • 1970-01-01
    • 2017-08-13
    • 1970-01-01
    相关资源
    最近更新 更多