【问题标题】:Queries with streaming sources must be executed with writeStream.start();必须使用 writeStream.start() 执行带有流源的查询;
【发布时间】:2017-03-29 07:57:20
【问题描述】:

我正在尝试在 spark 中读取来自 kafka(版本 10)的消息并尝试打印它。

     import spark.implicits._

         val spark = SparkSession
              .builder
              .appName("StructuredNetworkWordCount")
              .config("spark.master", "local")
              .getOrCreate()  

            val ds1 = spark.readStream.format("kafka")
              .option("kafka.bootstrap.servers", "localhost:9092")  
              .option("subscribe", "topicA")
              .load()

           ds1.collect.foreach(println)
           ds1.writeStream
           .format("console")
           .start()

           ds1.printSchema()

在线程“main”中得到一个错误异常

org.apache.spark.sql.AnalysisException:带有流式源的查询 必须用 writeStream.start();;;

执行

【问题讨论】:

    标签: scala apache-spark-sql spark-streaming


    【解决方案1】:

    您正在对查询计划进行分支:从您尝试的同一 ds1:

    • ds1.collect.foreach(...)
    • ds1.writeStream.format(...){...}

    但是您只在第二个分支上调用 .start(),而让另一个悬空而没有终止,这反过来又会引发您返回的异常。

    解决方案是启动两个分支并等待终止。

    val ds1 = spark.readStream.format("kafka")
      .option("kafka.bootstrap.servers", "localhost:9092")  
      .option("subscribe", "topicA")  
      .load()
    val query1 = ds1.collect.foreach(println)
      .writeStream
      .format("console")
      .start()
    val query2 = ds1.writeStream
      .format("console")
      .start()
    
    ds1.printSchema()
    query1.awaitTermination()
    query2.awaitTermination()
    

    【讨论】:

    • 那么解决方法是什么?
    • 我支持这里的评论。我们可以在这里得到一个合适的解决方案吗?也许是代码示例?谢谢!
    • 见@Rajeev 回答,awaitTermination 应该在start 之后调用
    【解决方案2】:

    我在这个问题上苦苦挣扎。我尝试了来自各种博客的每个建议的解决方案。 但我的情况是,在查询时调用 start() 和最后我调用 awaitTerminate() 函数之间几乎没有语句导致这种情况。

    请以这种方式尝试,它非常适合我。 工作示例:

    val query = df.writeStream
          .outputMode("append")
          .format("console")
          .start().awaitTermination();
    

    如果这样写会导致异常/错误:

    val query = df.writeStream
          .outputMode("append")
          .format("console")
          .start()
    
        // some statement 
        // some statement 
    
        query.awaitTermination();
    

    将抛出给定的异常并关闭您的流驱动程序。

    【讨论】:

    • 使用 Java 结构化流为我工作。我根本没有//some statements - 只是将StreamingQuery 保存到一个变量中,然后立即调用sQueryVar.start() 并遇到了同样的问题。这解决了它 - 谢谢!
    • 有谁知道这是什么原因?为什么 start() 和 awaitTermination() 之间的附加行会导致问题?
    【解决方案3】:

    我使用以下代码修复了问题。

     val df = session
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", brokers)
      .option("subscribe", "streamTest2")
      .load();
    
        val query = df.writeStream
      .outputMode("append")
      .format("console")
      .start()
    query.awaitTermination()
    

    【讨论】:

      【解决方案4】:

      请移除 ds1.collect.foreach(println)ds1.printSchema() ,将 outputModeawaitAnyTermination 用于后台进程 等待相关联的 spark.streams 上的任何查询终止

      val spark = SparkSession
          .builder
          .appName("StructuredNetworkWordCount")
          .config("spark.master", "local[*]")
          .getOrCreate()
      
        val ds1 = spark.readStream.format("kafka")
          .option("kafka.bootstrap.servers", "localhost:9092")
          .option("subscribe", "topicA")  .load()
      
        val consoleOutput1 = ds1.writeStream
           .outputMode("update")
           .format("console")
           .start()
      
        spark.streams.awaitAnyTermination()
      

      |key|value|topic|partition|offset|
      +---+-----+-----+---------+------+
      +---+-----+-----+---------+------+
      

      【讨论】:

      • 抱歉,只是出于好奇:为什么好心?您正在回答 OP 的问题,而不是要求 :)
      • 这里提到没有来自 kafka 的数据。请尝试从该主题发送数据
      【解决方案5】:

      我能够通过以下代码解决此问题。在我的场景中,我有多个中间数据帧,它们基本上是在 inputDF 上进行的转换。

       val query = joinedDF
            .writeStream
            .format("console")
            .option("truncate", "false")
            .outputMode(OutputMode.Complete())
            .start()
            .awaitTermination()
      

      joinedDF 是上次执行的转换的结果。

      【讨论】:

      • 如果你有一系列的转换,这很有效,但如果你有一个中间动作,你就完蛋了
      • @Mehdi LAMRANI 是的。这就是为什么我提到我使用这个选项的原因。
      • 我正在尝试解决后者,因为我偶然发现了这篇文章,因此我的评论:)
      猜你喜欢
      • 1970-01-01
      • 2021-10-13
      • 2021-01-31
      • 2018-03-14
      • 2021-08-16
      • 1970-01-01
      • 1970-01-01
      • 2017-06-23
      • 2019-05-25
      相关资源
      最近更新 更多