【发布时间】:2017-05-20 11:14:33
【问题描述】:
我有一个 (key, value) 的 RDD,我将其转换为 (key, List(value1, value2, value3) 的 RDD,如下所示。
val rddInit = sc.parallelize(List((1, 2), (1, 3), (2, 5), (2, 7), (3, 10)))
val rddReduced = rddInit..groupByKey.mapValues(_.toList)
rddReduced.take(3).foreach(println)
这段代码给了我下一个 RDD : (1,List(2, 3)) (2,List(5, 7)) (3,List(10))
但现在我想从我刚刚计算的 rdd(rddReduced rdd)返回到 rddInit。
我的第一个猜测是在键和 List 的每个元素之间实现某种叉积,如下所示:
rddReduced.map{
case (x, y) =>
val myList:ListBuffer[(Int, Int)] = ListBuffer()
for(element <- y) {
myList+=new Pair(x, element)
}
myList.toList
}.flatMap(x => x).take(5).foreach(println)
使用此代码,我得到了初始 RDD。但我不认为在火花作业中使用 ListBuffer 是一个好习惯。有没有其他方法可以解决这个问题?
【问题讨论】:
-
map后跟flatMap(identity)=>flatMap。使用element.map(Pair(...))- 使用ListBuffer会使代码过于复杂。将Pair设为案例类。
标签: scala apache-spark rdd