【发布时间】:2015-02-10 16:21:17
【问题描述】:
在http://spark.apache.org/docs/latest/programming-guide.html 阅读 Apache Spark 指南,它指出:
为什么 take 函数不能并行运行?并行实现此类功能有哪些困难?为了获取 RDD 的前 n 个元素,需要遍历整个 RDD,这是否与事实有关?
【问题讨论】:
标签: scala parallel-processing apache-spark
在http://spark.apache.org/docs/latest/programming-guide.html 阅读 Apache Spark 指南,它指出:
为什么 take 函数不能并行运行?并行实现此类功能有哪些困难?为了获取 RDD 的前 n 个元素,需要遍历整个 RDD,这是否与事实有关?
【问题讨论】:
标签: scala parallel-processing apache-spark
实际上,虽然take 不是完全并行的,但也不是完全顺序的。
例如,假设您take(200),每个分区有 10 个元素。 take 将首先获取分区 0 并查看它有 10 个元素。它假设需要 20 个这样的分区才能获得 200 个元素。但最好在并行请求中要求更多。所以它想要 30 个分区,而它已经有 1 个。所以它接下来会并行地获取分区 1 到 29。这很可能是最后一步。如果很不走运,一共没有找到200个元素,它会再次进行估计并并行请求另一个批次。
查看代码,有据可查: https://github.com/apache/spark/blob/v1.2.0/core/src/main/scala/org/apache/spark/rdd/RDD.scala#L1049
我认为文档是错误的。本地计算仅在需要单个分区时发生。这是第一次传递(获取分区 0)的情况,但通常不是后面传递的情况。
【讨论】:
take 将需要 2 次迭代,就像你描述的那样。
您将如何并行实现它?假设您有 4 个分区并且想要获取前 5 个元素。如果您事先知道每个分区的大小,这将很容易:例如,如果每个分区有 3 个元素,驱动程序会要求分区 0 获取所有元素,并要求分区 1 获取 2 个元素。所以问题是它不知道每个分区有多少元素。
现在,您可以先计算分区大小,但这需要限制支持的 RDD 转换集、多次计算元素或进行其他权衡,并且通常需要更多的通信开销。
【讨论】: