【问题标题】:How to reverse the result of reduceByKey using RDD API?如何使用 RDD API 反转 reduceByKey 的结果?
【发布时间】: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


【解决方案1】:

我很惊讶没有人提供具有 Scala 的 for-comprehension 的解决方案(在编译时“去糖”为 flatMapmap)。

我不经常使用这种语法,但是当我使用时……我觉得它很有趣。有些人更喜欢理解一系列flatMapmap,尤其是。用于更复杂的转换。

// that's what you ended up with after `groupByKey.mapValues`
val rddReduced: RDD[(Int, List[Int])] = ...
val r = for {
  (k, values) <- rddReduced
  v <- values
} yield (k, v)

scala> :type r
org.apache.spark.rdd.RDD[(Int, Int)]

scala> r.foreach(println)
(3,10)
(2,5)
(2,7)
(1,2)
(1,3)

// even nicer to our eyes
scala> r.toDF("key", "value").show
+---+-----+
|key|value|
+---+-----+
|  1|    2|
|  1|    3|
|  2|    5|
|  2|    7|
|  3|   10|
+---+-----+

毕竟,这就是我们享受 Scala 灵活性的原因,不是吗?

【讨论】:

  • 非常感谢@Jacked,我刚刚对你的答案投了赞成票。这个语法真的很神奇,我什至不知道它存在
【解决方案2】:

使用这种操作显然不是一个好习惯。

根据我在 spark-summit 课程中学到的知识,您必须尽可能多地使用Dataframes 和Datasets,使用它们您将受益于 spark 引擎的许多优化。

你想做的事情叫做explode,它是通过应用sql.functions包中的explode方法来实现的

解决方案应该是这样的:

 import spark.implicits._
 import org.apache.spark.sql.functions.explode
 import org.apache.spark.sql.functions.collect_list

 val dfInit = sc.parallelize(List((1, 2), (1, 3), (2, 5), (2, 7), (3, 10))).toDF("x", "y")
 val dfReduced = dfInit.groupBy("x").agg(collect_list("y") as "y")
 val dfResult = dfReduced.withColumn("y", explode($"y"))

dfResult 将包含与dfInit 相同的数据

【讨论】:

    【解决方案3】:

    这是将分组的 RDD 恢复为原始的一种方法:

    val rddRestored = rddReduced.flatMap{
        case (k, v) => v.map((k, _))
      }
    
    rddRestored.collect.foreach(println)
    (1,2)
    (1,3)
    (2,5)
    (2,7)
    (3,10)
    

    【讨论】:

      【解决方案4】:

      根据你的问题,我认为这就是你想要做的

      rddReduced.map{case(x, y) => y.map((x,_))}.flatMap(_).take(5).foreach(println)
      

      您会在分组后获得一个列表,您可以在其中再次映射它。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-06-05
        • 1970-01-01
        • 2011-09-27
        • 2015-12-09
        • 1970-01-01
        相关资源
        最近更新 更多