【问题标题】:Is there any better method than collect to read an RDD in spark?有没有比收集更好的方法来读取 Spark 中的 RDD?
【发布时间】:2015-07-31 18:18:06
【问题描述】:

所以,我想将 RDD 读入一个数组。为此,我可以使用 collect 方法。但是这种方法真的很烦人,因为在我的情况下它不断给出 kyro 缓冲区溢出错误。如果我将 kyro 缓冲区大小设置得太大,它就会开始出现问题。另一方面,我注意到如果我只是使用 saveAsTextFile 方法将 RDD 保存到文件中,我不会收到任何错误。所以,我在想,一定有更好的方法将 RDD 读入数组,它不像 collect 方法那样有问题。

【问题讨论】:

  • 也许最好问一个关于缓冲区溢出错误的问题。
  • 我可以设置的 Kyro Serializer 缓冲区的最大大小为 1GB。那么,这是否意味着我不能收集超过 1GB 的数据?
  • 也许切换到标准序列化?
  • @maasg:你如何切换到标准序列化?
  • 一个更相关的问题是 a) 为什么不能保存文件并读回文件? b) 为什么你真的希望在一个驱动程序中拥有 >1G 的数据?

标签: java serialization apache-spark bigdata


【解决方案1】:

没有。 collect 是将 RDD 读入数组的唯一方法。

saveAsTextFile 永远不需要将所有数据收集到一台机器上,因此它不会像collect 那样受到单台机器上可用内存的限制。

【讨论】:

  • 但是为什么这个 collect 方法总是抛出这些缓冲区溢出错误。此外,我无法将 Kyro 序列化程序的最大大小设置为超过 1GB。那么,似乎不可能收集大于 1GB 的东西?知道如何解决这个问题吗?
  • 我不知道你看到了什么错误。也不知道为什么不能将设置(哪个设置?)​​增加到 1GB 以上。 (您是否也增加了spark.driver.maxResultSize?)无论如何,您可以尝试的另一件事是使用RDD.toLocalIterator 而不是collect。只是在黑暗中拍摄。
  • 还有一个每个分区 2GB 的限制,您可能即将达到。 (issues.apache.org/jira/browse/SPARK-4885) 在这种情况下,更多的分区是解决方案。
【解决方案2】:

toLocalIterator()

此方法返回一个包含此 RDD 中所有元素的迭代器。迭代器将消耗与此 RDD 中最大分区一样多的内存。作为 RunJob 处理以在每个步骤上评估一个单独的分区。

>>> x = rdd.toLocalIterator()
>>> x
<generator object toLocalIterator at 0x283cf00>

那么你可以通过

访问rdd中的元素
empty_array = []    
for each_element in x:
    empty_array.append(each_element)

https://spark.apache.org/docs/1.0.2/api/java/org/apache/spark/rdd/RDD.html#toLocalIterator()

【讨论】:

    猜你喜欢
    • 2012-03-15
    • 2011-07-20
    • 2015-09-10
    • 2019-10-17
    • 2011-03-09
    • 2011-10-19
    • 2012-03-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多