【发布时间】: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