【问题标题】:map vs mapValues in SparkSpark中的map与mapValues
【发布时间】:2016-08-10 07:54:30
【问题描述】:

我目前正在学习 Spark 并开发自定义机器学习算法。我的问题是.map().mapValues() 之间有什么区别,在哪些情况下我显然必须使用其中一个而不是另一个?

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    mapValues 仅适用于 PairRDD,即 RDD[(A, B)] 形式的 RDD。在这种情况下,mapValues 仅对 value (元组的第二部分)进行操作,而 map整个记录 (键和值的元组)上操作)。

    换句话说,给定f: B => Crdd: RDD[(A, B)],这两个是相同的(几乎 - 见底部的评论):

    val result: RDD[(A, C)] = rdd.map { case (k, v) => (k, f(v)) }
    
    val result: RDD[(A, C)] = rdd.mapValues(f)
    

    后者更短更清晰,所以当你只想转换值并保持键不变时,建议使用mapValues

    另一方面,如果你也想转换键(例如你想应用f: (A, B) => C),你根本不能使用mapValues,因为它只会将值传递给你的函数。

    最后一个区别在于分区:如果你对你的RDD应用了任何自定义分区(例如使用partitionBy),使用map会“忘记”那个分区器(结果将恢复为默认值)分区),因为键可能已更改;但是,mapValues 会保留 RDD 上设置的任何分区器。

    【讨论】:

    • 我想知道它们是否对性能有影响,因为我正在尝试优化this,但从你所说的来看,我想它不会有任何影响......
    • @gsamaras 它可能会影响性能,因为如果您需要使用相同的密钥再次重新分区,丢失分区信息将迫使您重新分区。
    【解决方案2】:

    当我们使用带有 Pair RDD 的 map() 时,我们可以访问 Key 和 value。有几次我们只对访问值(而不是键)感兴趣。在这种情况下,我们可以使用 mapValues() 而不是 map()。

    ma​​pValues 示例

    val inputrdd = sc.parallelize(Seq(("maths", 50), ("maths", 60), ("english", 65)))
    val mapped = inputrdd.mapValues(mark => (mark, 1));
    
    //
    val reduced = mapped.reduceByKey((x, y) => (x._1 + y._1, x._2 + y._2))
    
    reduced.collect
    

    Array[(String, (Int, Int))] = Array((english,(65,1)), (maths,(110,2)))

    val average = reduced.map { x =>
                               val temp = x._2
                               val total = temp._1
                               val count = temp._2
                               (x._1, total / count)
                               }
    
    average.collect()
    

    res1: Array[(String, Int)] = Array((english,65), (maths,55))

    【讨论】:

      【解决方案3】:

      map 采用一个函数来转换集合的每个元素:

       map(f: T => U)
      RDD[T] => RDD[U]
      

      T 是一个元组时,我们可能只想作用于值而不是键 mapValues 采用一个函数,将输入中的值映射到输出中的值:mapValues(f: V => W) 在哪里RDD[ (K, V) ] => RDD[ (K, W) ]

      提示:使用mapValues可以避免reshuffle当数据按key进行分区时

      【讨论】:

        【解决方案4】:
        val inputrdd = sc.parallelize(Seq(("india", 250), ("england", 260), ("england", 180)))
        

        (1)

        map():-
        
        val mapresult= inputrdd.map{b=> (b,1)}
        mapresult.collect
        
        Result-= Array(((india,250),1), ((england,260),1), ((english,180),1))
        

        (2)

        mapvalues():-
        
        val mapValuesResult= inputrdd.mapValues(b => (b, 1));
        mapValuesResult.collect
        

        结果-

        Array((india,(250,1)), (england,(260,1)), (england,(180,1)))
        

        【讨论】:

        • 除非您完全确定可以将所有数据放入内存中,否则永远不要尝试使用.collect
        猜你喜欢
        • 1970-01-01
        • 2014-10-27
        • 1970-01-01
        • 2018-06-10
        • 2015-03-16
        • 2016-07-18
        • 1970-01-01
        • 2020-03-07
        • 1970-01-01
        相关资源
        最近更新 更多