【问题标题】:multiple writeStream with spark streaming带有火花流的多个 writeStream
【发布时间】:2020-02-24 21:13:06
【问题描述】:

我正在使用 spark 流,但在尝试实现多个 writestream 时遇到了一些问题。 下面是我的代码

DataWriter.writeStreamer(firstTableData,"parquet",CheckPointConf.firstCheckPoint,OutputConf.firstDataOutput)
DataWriter.writeStreamer(secondTableData,"parquet",CheckPointConf.secondCheckPoint,OutputConf.secondDataOutput)
DataWriter.writeStreamer(thirdTableData,"parquet", CheckPointConf.thirdCheckPoint,OutputConf.thirdDataOutput)

其中 writeStreamer 定义如下:

def writeStreamer(input: DataFrame, checkPointFolder: String, output: String) = {

  val query = input
                .writeStream
                .format("orc")
                .option("checkpointLocation", checkPointFolder)
                .option("path", output)
                .outputMode(OutputMode.Append)
                .start()

  query.awaitTermination()
}

我面临的问题是只有第一个表是用 spark writeStream 编写的,所有其他表都没有任何反应。 请问您对此有什么想法吗?

【问题讨论】:

    标签: apache-spark spark-structured-streaming


    【解决方案1】:

    query.awaitTermination() 应该在创建最后一个流之后完成。

    writeStreamer 函数可以修改为返回 StreamingQuery 而不是此时的 awaitTermination(因为它是阻塞):

    def writeStreamer(input: DataFrame, checkPointFolder: String, output: String): StreamingQuery = {
      input
        .writeStream
        .format("orc")
        .option("checkpointLocation", checkPointFolder)
        .option("path", output)
        .outputMode(OutputMode.Append)
        .start()
    }
    

    那么你将拥有:

    val query1 = DataWriter.writeStreamer(...)
    val query2 = DataWriter.writeStreamer(...)
    val query3 = DataWriter.writeStreamer(...)
    
    query3.awaitTermination()
    

    【讨论】:

      【解决方案2】:

      如果您想执行编写器以并行运行,您可以使用

      sparkSession.streams.awaitAnyTermination()
      

      并从 writeStreamer 方法中删除 query.awaitTermination()

      【讨论】:

        【解决方案3】:

        默认情况下,并发作业的数量是 1,这意味着一次 只有 1 个工作将处于活动状态

        您是否尝试在 spark conf 中增加可能的并发作业数量?

        sparkConf.set("spark.streaming.concurrentJobs","3")
        

        非官方来源:http://why-not-learn-something.blogspot.com/2016/06/spark-streaming-performance-tuning-on.html

        【讨论】:

        • 我从 writer 中删除了 await 终止,现在我收到一个新错误:ERROR org.apache.spark.sql.execution.datasources.FileFormatWriter: Aborting job null。 org.apache.spark.SparkException: Job 1 由于 SparkContext 已关闭而取消
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-04-27
        • 1970-01-01
        • 2018-12-17
        • 2018-01-08
        相关资源
        最近更新 更多