【问题标题】:The streamingContext stops before waiting for the processing of all received data to be completedstreamingContext 在等待所有接收到的数据处理完成之前停止
【发布时间】:2018-01-26 06:56:36
【问题描述】:

我想第一次使用和了解 Spark Streaming。我测试了这个例子:

val queueOfRDDs:Queue[RDD[Int]] = Queue.empty[RDD[Int]]        
@transient val streamingContext:StreamingContext = new StreamingContext(sc, Seconds(1))
val inputDStream = streamingContext.queueStream(queueOfRDDs,true,null)
inputDStream.foreachRDD(rdd =>
{
    if(!rdd.isEmpty())
        println("size of rdd "+rdd.count())
    else
        {
        println("empty rdd")        
        }

})
streamingContext.start()
queueOfRDDs.synchronized {
  for(a <- 1 to 10)
  {
    queueOfRDDs.+=(Config.sc.makeRDD(1 to 1000, 10))        
  }
}
streamingContext.stop(false,true)

我得到:

empty rdd
size of rdd 1000

我放了“stopGracefully = true”(streamingContext.stop(false,true)),因为我想处理队列中的所有rdds,而streamingContext在等待所有接收到的数据处理完成后停止。但是,只有一个 rdd 被处理。 你能帮帮我吗

【问题讨论】:

  • 你怎么知道只处理了一个 RDD,请附上你得到的结果和你真正想要的
  • 我想处理队列中的所有 rdds。所以,预期的结果是: rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 大小 rdd 1000 但我得到这个结果: rdd 1000 的空 rdd 大小
  • 预期结果是处理了10个RDD但只处理了一个

标签: apache-spark streaming rdd


【解决方案1】:

火花流旨在: 1、用StreamingQuery.awaitTermination()永远运行:Unit 2、使用StreamingQuery.awaitTermination(timeoutMs: Long): Boolean 运行直到超时

流不知道数据是否结束。你调用streamingContext.stop,程序直接结束。 在您的情况下,您可以使用 awaitTermination(timeoutMs: Long)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-02-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-27
    • 1970-01-01
    相关资源
    最近更新 更多