【问题标题】:in scala,how to get the collect the values in Iterable[(Tuple)]在 Scala 中,如何获取 Iterable[(Tuple)] 中的值
【发布时间】:2016-06-11 06:18:05
【问题描述】:

我有一个RDD[(Key,Iterable[(Name,Value)])]。 我正在尝试获取Name 中的所有值,以便我可以获取Name 的所有唯一出现,然后创建一个Index,以便我可以创建RDD[(Key,Iterable[(Index,Value)])] 的结果RDD

输入示例:

(4048,CompactBuffer(("a",3.0), ("b",9.0), ("c",14.0))
(4049,CompactBuffer(("a",2.0), ("c",14.0))
(4050,CompactBuffer(("b",2.0), ("d",10.0))

输出示例:

(4048,CompactBuffer((1,3.0), (2,9.0), (3,14.0))
(4049,CompactBuffer((1,2.0), (3,12.0))
(4050,CompactBuffer((2,2.0), (4,10.0))

【问题讨论】:

    标签: scala rdd iterable


    【解决方案1】:

    如果我正确理解你想要的东西,这样的事情应该可以工作:

    val rdd: RDD[K, Iterable[N, V]] = ???
    
    val nameIndexMap = rdd.flatMap(_._2.map(nv => nv._1))
                          .distinct
                          .collect
                          .zipWithIndex
                          .toMap
    
    val newRDD = rdd.mapValues(xs => xs.map(nv => (nameIndexMap(nv._1), nv._2)))
    

    如果nameIndexMap 很大或将被多次使用,那么您可能想要使用广播,例如像这样的

    val nameIndexMapBc = sc.broadcast(nameIndexMap)  //sc is the SparkContext
    
    val newRDD = rdd.mapValues(xs => xs.map(nv => (nameIndexMap.value)(nv._1), nv._2)))
    

    【讨论】:

      【解决方案2】:

      对于 jasonl 的一个小补充,如果名称空间的基数很高,您可以执行类似的转换,避免将名称收集到驱动程序内存中:

      val named = rdd.flatMap { case (key, values) => values.map { case (name, value) => (name, (key, value)) } }
      
      val nameMap = named.map(_._1).distinct().zipWithUniqueId()
      
      val indexed = named.join(nameMap).map { case (name, ((key, value), index)) => (key, (index, value)) }.groupByKey
      

      【讨论】:

        猜你喜欢
        • 2016-01-21
        • 1970-01-01
        • 2013-04-18
        • 2010-11-07
        • 2019-06-25
        • 1970-01-01
        • 1970-01-01
        • 2012-07-20
        • 2010-11-06
        相关资源
        最近更新 更多