【问题标题】:How to properly apply HashPartitioner before a join in Spark?如何在加入 Spark 之前正确应用 HashPartitioner?
【发布时间】:2019-08-12 06:37:32
【问题描述】:

为了减少两个 RDD 加入过程中的洗牌,我决定首先使用 HashPartitioner 对它们进行分区。这是我的做法。我这样做是正确的,还是有更好的方法来做到这一点?

val rddA = ...
val rddB = ...

val numOfPartitions = rddA.getNumPartitions

val rddApartitioned = rddA.partitionBy(new HashPartitioner(numOfPartitions))
val rddBpartitioned = rddB.partitionBy(new HashPartitioner(numOfPartitions))

val rddAB = rddApartitioned.join(rddBpartitioned)

【问题讨论】:

    标签: scala apache-spark rdd partitioner


    【解决方案1】:

    为了减少两个 RDD 连接过程中的洗牌,

    令人惊讶的普遍误解是,重新分区会减少甚至消除随机播放。 它没有。重新分区 是 shuffle 的最纯粹形式。它不会节省时间、带宽或内存。

    使用主动分区器背后的原理是不同的 - 它允许您洗牌一次,并重用状态,执行多个按键操作,而无需额外的洗牌(尽管据我所知,不一定没有额外的网络流量,as co-partitioning doesn't imply co-location,不包括在单个操作中发生随机播放的情况)。

    所以你的代码是正确的,但如果你加入一次它不会给你带来任何东西。

    【讨论】:

    • 很好的观察 :) 就我而言,之后我也做了一个 sortByKey,所以我想它会有所帮助。
    • 如果sortByKey 应用在rddAB 上,则完全没有区别。如果将其应用于rddApartitioned / rddBpartitioned 那么它可以提供一些好处。
    • @user10938362 是否要求两个 RDD 都是 partitionByed?如果我只为rddA 而不是rddB 提供partitionBy 会发生什么?
    【解决方案2】:

    只有一条评论,如果rddApartitionedrddBpartitioned有多个动作,最好在.partitionBy之后附加.persist(),否则,所有动作都会评估rddApartitionedrddBpartitioned的整个谱系,这将导致哈希分区一次又一次地发生。

    val rddApartitioned = rddA.partitionBy(new HashPartitioner(numOfPartitions)).persist()
    val rddBpartitioned = rddB.partitionBy(new HashPartitioner(numOfPartitions)).persist()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-07-28
      相关资源
      最近更新 更多