【问题标题】:In Apache Spark cogroup, how to make sure 1 RDD of >2 operands is not moved?在 Apache Spark cogroup 中,如何确保不移动 >2 个操作数的 1 个 RDD?
【发布时间】:2017-12-14 09:25:19
【问题描述】:

在同群变换中,例如RDD1.cogroup(RDD2, ...),我曾经假设 Spark 仅在以下情况下洗牌/移动 RDD2 并保留 RDD1 的分区和内存存储:

  1. RDD1 有一个显式分区器
  2. RDD1 被缓存。

在我的其他项目中,大多数洗牌行为似乎与这个假设一致。所以昨天我写了一个简短的 scala 程序来一劳永逸地证明它:

// sc is the SparkContext
val rdd1 = sc.parallelize(1 to 10, 4).map(v => v->v)
  .partitionBy(new HashPartitioner(4))
rdd1.persist().count()
val rdd2 = sc.parallelize(1 to 10, 4).map(v => (11-v)->v)

val cogrouped = rdd1.cogroup(rdd2).map {
  v =>
    v._2._1.head -> v._2._2.head
}

val zipped = cogrouped.zipPartitions(rdd1, rdd2) {
  (itr1, itr2, itr3) =>
    itr1.zipAll(itr2.map(_._2), 0->0, 0).zipAll(itr3.map(_._2), (0->0)->0, 0)
      .map {
        v =>
          (v._1._1._1, v._1._1._2, v._1._2, v._2)
      }
}

zipped.collect().foreach(println)

如果 rdd1 不移动 zipped 的第一列应该和第三列的值相同,所以我运行了程序,哎呀:

(4,7,4,1)
(8,3,8,2)
(1,10,1,3)
(9,2,5,4)
(5,6,9,5)
(6,5,2,6)
(10,1,6,7)
(2,9,10,0)
(3,8,3,8)
(7,4,7,9)
(0,0,0,10)

假设不正确。 Spark 可能进行了一些内部优化,并决定重新生成 rdd1 的分区比将它们保存在缓存中要快得多。

所以问题是:如果我不移动 RDD1(并保持缓存)的编程要求是由于速度以外的其他原因(例如资源局部性),或者在某些情况下 Spark 内部优化不是可取的,有没有办法明确指示框架不要在所有类似 cogroup 的操作中移动操作数?这还包括联接、外部联接和 groupWith。

非常感谢您的帮助。到目前为止,我使用广播连接作为一种不可扩展的临时解决方案,它不会持续很长时间,然后我的集群就会崩溃。我期待一个与分布式计算主体一致的解决方案。

【问题讨论】:

    标签: scala apache-spark join shuffle


    【解决方案1】:

    如果 rdd1 不移动 zipped 的第一列应该和第三列有相同的值

    这个假设是不正确的。创建CoGroupedRDD不仅是shuffle,还包括生成匹配对应记录所需的内部结构。在内部,Spark 将使用自己的 ExternalAppendOnlyMap,它使用自定义的开放哈希表实现 (AppendOnlyMap),它不提供任何排序​​保证。

    如果你检查调试字符串:

    zipped.toDebugString
    
    (4) ZippedPartitionsRDD3[8] at zipPartitions at <console>:36 []
     |  MapPartitionsRDD[7] at map at <console>:31 []
     |  MapPartitionsRDD[6] at cogroup at <console>:31 []
     |  CoGroupedRDD[5] at cogroup at <console>:31 []
     |  ShuffledRDD[2] at partitionBy at <console>:27 []
     |      CachedPartitions: 4; MemorySize: 512.0 B; ExternalBlockStoreSize: 0.0 B; DiskSize: 0.0 B
     +-(4) MapPartitionsRDD[1] at map at <console>:26 []
        |  ParallelCollectionRDD[0] at parallelize at <console>:26 []
     +-(4) MapPartitionsRDD[4] at map at <console>:29 []
        |  ParallelCollectionRDD[3] at parallelize at <console>:29 []
     |  ShuffledRDD[2] at partitionBy at <console>:27 []
     |      CachedPartitions: 4; MemorySize: 512.0 B; ExternalBlockStoreSize: 0.0 B; DiskSize: 0.0 B
     +-(4) MapPartitionsRDD[1]...
    

    您会看到 Spark 确实使用 CachedPartitions 来计算 zipped RDD。如果您还跳过删除分区器的map 转换,您会看到coGroup 重用了rdd1 提供的分区器:

    rdd1.cogroup(rdd2).partitioner == rdd1.partitioner
    
    Boolean = true
    

    【讨论】:

      猜你喜欢
      • 2016-04-26
      • 1970-01-01
      • 1970-01-01
      • 2015-06-15
      • 2019-02-11
      • 1970-01-01
      • 2015-11-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多