我想在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.有没有相同的方法来做这样的事情?
其实我完全误解了由多个RDD组成的DStream的概念。一个DStream是由多个RDD组成的,但是随着时间的推移。
在微批次的上下文中,DStream 由当前的 RDD 组成。
所以,上面的代码总是有效的。
| 归档时间: |
|
| 查看次数: |
888 次 |
| 最近记录: |