【问题标题】:Can only zip RDDs with same number of elements in each partition despite repartition尽管重新分区,但只能压缩每个分区中元素数量相同的 RDD
【发布时间】:2016-07-03 13:31:41
【问题描述】:

我加载一个数据集

val data = sc.textFile("/home/kybe/Documents/datasets/img.csv",defp)

我想为此数据建立一个索引

val nb = data.count.toInt
val tozip = sc.parallelize(1 to nb).repartition(data.getNumPartitions)

val res = tozip.zip(data)

不幸的是我有以下错误

Can only zip RDDs with same number of elements in each partition

如果可能的话,如何按分区修改元素的数量?

【问题讨论】:

  • 我是一名初学者,所以可能不是最好的解决方案,但理论上可以在两个 RDD 上 zipWithIndex 然后执行 join (或左连接或右连接,具体取决于你想如何混合两个RDD)使用元素的索引,另一种方法是计算长度差异并用默认值填充空白,但不确定如何使用RDD。
  • 我完全忘记了这个伟大的功能!谢谢,但我还是不明白为什么它不起作用。
  • AFAIK 我认为拥有相同数量的分区并不等于拥有两个长度相同的 RDD,我认为 repartition 会从文档中对集群中的分区进行洗牌:Reshuffle the data in the RDD randomly to create either more or fewer partitions and balance it across them. This always shuffles all data over the network. ,如果我做对了整个事情,那么在你的情况下,如果你需要加入(这会导致另一次改组),那么你的情况就没有多大意义了。

标签: scala apache-spark rdd


【解决方案1】:

为什么它不起作用?

zip() 的文档指出:

用另一个 RDD 压缩这个 RDD,返回每个 RDD 中的第一个元素、每个 RDD 中的第二个元素等的键值对。假设两个 RDD 具有相同数量的分区每个分区中的元素数量相同(例如,一个是通过另一个映射生成的)。

所以我们需要确保我们满足两个条件:

  • 两个 RDD 的分区数相同
  • 这些 RDD 中的各个分区具有完全相同的大小

您确保您将拥有与repartition() 相同数量的分区,但 Spark 不保证您在每个 RDD 的每个分区中拥有相同的分布。

为什么会这样?

因为有不同类型的 RDD,而且大多数都有不同的分区策略!例如:

  • ParallelCollectionRDD 是在您使用sc.parallelize(collection) 并行化集合时创建的,它将查看应该有多少个分区,将检查集合的大小并计算step 的大小。 IE。您在列表中有 15 个元素并想要 4 个分区,前 3 个将有 4 个连续的元素,最后一个将有剩余的 3 个。
  • HadoopRDD,如果我没记错的话,每个文件块一个分区。即使您在内部使用本地文件,Spark 在读取本地文件时首先创建这种 RDD,然后映射该 RDD,因为该 RDD 是 <Long, Text> 的一对 RDD,而您只想要 String :-)
  • 等等等等

在您的示例中,Spark 内部在进行重新分区时确实创建了不同类型的 RDD(CoalescedRDDShuffledRDD),但我认为您有一个全局想法,即 不同的 RDD 具有不同的分区策略: -)

注意zip() 文档的最后一部分提到了map() 操作。此操作不会重新分区,因为它是一个转换数据,因此它将保证这两个条件。

解决方案

在上面提到的这个简单示例中,您可以简单地使用data.zipWithIndex。如果您需要更复杂的东西,那么为zip() 创建新的RDD 应该使用map() 创建,如上所述。

【讨论】:

    【解决方案2】:

    我通过创建一个像这样的隐式助手解决了这个问题

    implicit class RichContext[T](rdd: RDD[T]) {
      def zipShuffle[A](other: RDD[A])(implicit kt: ClassTag[T], vt: ClassTag[A]): RDD[(T, A)] = {
        val otherKeyd: RDD[(Long, A)] = other.zipWithIndex().map { case (n, i) => i -> n }
        val thisKeyed: RDD[(Long, T)] = rdd.zipWithIndex().map { case (n, i) => i -> n }
        val joined                    = new PairRDDFunctions(thisKeyed).join(otherKeyd).map(_._2)
        joined
      }
    }
    

    然后可以像这样使用

    val rdd1   = sc.parallelize(Seq(1,2,3))
    val rdd2   = sc.parallelize(Seq(2,4,6))
    val zipped = rdd1.zipShuffle(rdd2) // Seq((1,2),(2,4),(3,6))
    

    注意:请记住,join 会导致随机播放。

    【讨论】:

    • Join 是 O(nm),集合的 n 和 m 大小与 zip 相比是 O(n),因此 join 对于大数据量是不够的。
    【解决方案3】:

    下面通过定义 custom_zip 方法为这个问题提供了 Python 答案: Can only zip with RDD which has the same number of partitions error

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-06-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多