【发布时间】:2017-12-14 09:25:19
【问题描述】:
在同群变换中,例如RDD1.cogroup(RDD2, ...),我曾经假设 Spark 仅在以下情况下洗牌/移动 RDD2 并保留 RDD1 的分区和内存存储:
- RDD1 有一个显式分区器
- 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