【问题标题】:Is groupByKey ever preferred over reduceByKeygroupByKey 是否比 reduceByKey 更受欢迎
【发布时间】:2016-01-18 05:32:24
【问题描述】:

当我需要在 RDD 中对数据进行分组时,我总是使用reduceByKey,因为它在打乱数据之前会执行 map 侧归约,这通常意味着更少的数据被打乱,从而获得更好的性能。即使map端reduce函数收集了所有的值并没有真正减少数据量,我仍然使用reduceByKey,因为我假设reduceByKey的性能永远不会比groupByKey差。但是,我想知道这个假设是否正确,或者是否确实存在应该首选groupByKey 的情况??

【问题讨论】:

  • 根据我在下面得到的答案(感谢您的回答),@eliasah 说groupByKey 只是语法糖,而@climbage 认为reduceByKey 如果我使用它可能会稍微慢一些复制groupByKey 功能。我想我实际上会尝试在一些示例中测试这两个函数:)
  • 我需要使用 groupByKey 的唯一一次是对依赖于先前值的数据样本进行计算。一个预先计算的运行总计就是一个例子。 GPS距离。等等。

标签: apache-spark rdd


【解决方案1】:

我不会发明轮子,根据代码文档,groupByKey 操作将 RDD 中每个键的值分组到一个序列中,这还允许通过以下方式控制结果键值对 RDD 的分区传递一个Partitioner。

此操作可能非常昂贵。如果您为了对每个键执行聚合(例如求和或平均)进行分组,则使用 aggregateByKey 或 reduceByKey 将提供更好的性能。

注意:按照目前的实现,groupByKey 必须能够在内存中保存任何键的所有键值对。如果键的值过多,可能会导致 OOME。

事实上,我更喜欢combineByKey 操作,但是如果你不是很熟悉map-reduce 范式,有时很难理解combiner 和merge 的概念。为此,您可以阅读 yahoo map-reduce 圣经here,它很好地解释了这个主题。

有关更多信息,我建议您阅读PairRDDFunctions code。

【讨论】:

  • 我了解与groupByKey 相关的可能问题(例如给定键的值太多) - 问题是有时groupByKey 实际上是更好的选择。您提到使用groupByKey时可以控制生成的键值对的分区,但也可以使用reduceByKey控制,所以这似乎不是使用groupByKey的理由,或者我我误会你了?
  • 完全正确,您可以将groupByKey 视为语法糖。如果可以避免,最好使用 aggregateByKey、reduceByKey 或 combineByKey
  • @GlennieHellesSindholt 你似乎不太相信。
  • combineByKey 如何避免 OOM 问题?结果大小相同。
  • @shuaiyuancn combineByKey 与CompactBuffer、+= 和++= 完全等价于groupByKey,combineByKey 允许您根据数据选择更高效的数据结构分配。可以说,只有少数情况无法通过重新分区和外部排序来代替,但对于典型用户来说,这很可能是低级方法。
【解决方案2】:

reduceByKey 和 groupByKey 都使用 combineByKey 和不同的组合/合并语义。

我看到的关键区别是groupByKey 将标志 (mapSideCombine=false) 传递给了随机播放引擎。从问题SPARK-772 来看,这是对shuffle 引擎的提示,在数据大小不会改变时不要运行mapside 组合器。

所以我想说,如果您尝试使用 reduceByKey 来复制 groupByKey,您可能会看到轻微的性能损失。

【讨论】:

    【解决方案3】:

    我相信climbage和eliasah忽略了问题的其他方面:

    • 代码可读性
    • 代码可维护性
    • 代码库大小

    如果操作不减少数据量,它必须以一种或另一种语义等效于GroupByKey。假设我们有RDD[(Int,String)]:

    import scala.util.Random
    Random.setSeed(1)
    
    def randomString = Random.alphanumeric.take(Random.nextInt(10)).mkString("")
    
    val rdd = sc.parallelize((1 to 20).map(_ => (Random.nextInt(5), randomString)))
    

    我们想要连接给定键的所有字符串。使用groupByKey 非常简单:

    rdd.groupByKey.mapValues(_.mkString(""))
    

    reduceByKey 的简单解决方案如下所示:

    rdd.reduceByKey(_ + _)
    

    它很短,可以说很容易理解,但存在两个问题:

    • 效率极低,因为它每次都会创建一个新的String 对象*
    • 建议您执行的操作比实际操作更便宜,尤其是在您只分析 DAG 或调试字符串时

    为了处理第一个问题,我们需要一个可变数据结构:

    import scala.collection.mutable.StringBuilder
    
    rdd.combineByKey[StringBuilder](
        (s: String) => new StringBuilder(s),
        (sb: StringBuilder, s: String) => sb ++= s,
        (sb1: StringBuilder, sb2: StringBuilder) => sb1.append(sb2)
    ).mapValues(_.toString)
    

    它仍然暗示了其他一些正在发生的事情并且非常冗长,特别是如果在您的脚本中重复多次。你当然可以提取匿名函数

    val createStringCombiner = (s: String) => new StringBuilder(s)
    val mergeStringValue = (sb: StringBuilder, s: String) => sb ++= s
    val mergeStringCombiners = (sb1: StringBuilder, sb2: StringBuilder) => 
      sb1.append(sb2)
    
    rdd.combineByKey(createStringCombiner, mergeStringValue, mergeStringCombiners)
    

    但归根结底,这仍然意味着要付出额外的努力来理解这段代码,增加复杂性并且没有真正的附加值。我发现特别令人不安的一件事是明确包含可变数据结构。即使 Spark 处理了几乎所有的复杂性,这也意味着我们不再拥有优雅的、引用透明的代码。

    我的观点是,如果您真的想尽一切办法减少数据量,请使用reduceByKey。否则,您的代码会更难编写、更难分析,而且一无所获。

    注意:

    此答案主要针对 Scala RDD API。当前的 Python 实现与其对应的 JVM 有很大不同,并且包括一些优化,在类似 groupBy 的操作的情况下,这些优化比简单的 reduceByKey 实现具有显着优势。

    对于Dataset API,请参阅DataFrame / Dataset groupBy behaviour/optimization。


    * 请参阅Spark performance for Scala vs Python 以获得令人信服的示例

    【讨论】:

      猜你喜欢
      • 2021-06-06
      • 2011-03-08
      • 2019-06-11
      • 2020-10-25
      • 2012-05-13
      • 2012-09-16
      • 2012-11-16
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多