如何计算前k个单词

Nic*_*aro 11 spark-streaming

我想在Spark-Streaming应用程序中计算前k个单词,在时间窗口中收集文本行.

我最终得到以下代码:

...
val window = stream.window(Seconds(30))

val wc = window
  .flatMap(line => line.split(" "))
  .map(w => (w, 1))
  .reduceByKey(_ + _)

wc.foreachRDD(rdd => {
  println("---------------------------------------------------")
  rdd.top(10)(Ordering.by(_._2)).zipWithIndex.foreach(println)
})
...
Run Code Online (Sandbox Code Playgroud)

它似乎工作.

问题:前k字图表是使用在(变量)返回的foreachRDD每个上执行top + print函数的函数计算的.RDDreduceByKeywc

事实证明,reduceByKey返回一个DStream单一的RDD,所以上面的代码工作,但规格不保证正确的行为.

我错了,它适用于所有情况吗?

为什么在spark-streaming中没有一种方法可以将a DStream作为单个RDD而不是RDD对象集合来考虑,以便执行更复杂的转换?

我的意思是这样的功能:dstream.withUnionRDD(rdd => ...)允许您对单个/联合进行转换和操作RDD.有没有相同的方法来做这样的事情?

Nic*_*aro 2

其实我完全误解了由多个RDD组成的DStream的概念。一个DStream是由多个RDD组成的,但是随着时间的推移。

在微批次的上下文中,DStream 由当前的 RDD 组成。

所以,上面的代码总是有效的。