1

对我从互联网上借来的用于研究目的的一段代码感到困惑。这是代码:

import org.apache.spark.sql.SparkSession
import org.apache.spark.rdd.RDD
import org.apache.spark.streaming.{Seconds, StreamingContext}
import scala.collection.mutable

val spark = ... 

val sc = spark.sparkContext
val ssc = new StreamingContext(spark.sparkContext, Seconds(1)) 

val rddQueue = new mutable.Queue[RDD[Char]]()
val QS = ssc.queueStream(rddQueue) 

QS.foreachRDD(q=> {
   print("Hello") // Queue never exhausted
   if(!q.isEmpty) {
       ... do something
       ... do something
   }
}
)

//ssc.checkpoint("/chkpoint/dir") if unchecked causes Serialization error

ssc.start()
for (c <- 'a' to 'c') {
    rddQueue += ssc.sparkContext.parallelize(List(c))
}
ssc.awaitTermination()

我正在追踪它只是为了检查并注意到“你好”被永远打印出来:

 HelloHelloHelloHelloHelloHelloHelloHelloHelloHello and so on

我原以为 queueStream 会在 3 次迭代后耗尽。

那么,我错过了什么?

4

1 回答 1

0

知道了。它实际上已经用尽了,但循环仍在继续,这就是为什么声明

 if(!q.isEmpty)

在那儿。

好的,本来以为它会停止,或者说不执行,但不是这样。我想起来了。如果没有流式传输,将根据批处理间隔的时间产生一个空的 RDD。为其他人留下,因为有一个赞成票。

然而,即使是遗留的,它也是一个不好的例子,因为添加检查点会导致序列化错误。为了他人的利益而离开它。

ssc.checkpoint("/chkpoint/dir")
于 2019-01-04T11:27:56.650 回答