【发布时间】:2021-10-12 10:25:46
【问题描述】:
我有一个 pyspark RDD,它有大约 200 万个元素。我不能一次收集它们,因为它会导致OutOfMemoryError 异常。
如何分批收集?
这是一个潜在的解决方案,但我怀疑有更好的解决方案:收集一个批次(使用 take、https://spark.apache.org/docs/3.1.2/api/python/reference/api/pyspark.RDD.take.html#pyspark.RDD.take),然后从该批次中的 RDD 中删除所有元素(使用 filter、https://spark.apache.org/docs/3.1.2/api/python/reference/api/pyspark.RDD.filter.html#pyspark.RDD.filter、但我怀疑有更好的方法),重复直到没有元素被收集。
【问题讨论】:
-
.collect()不打算在实验和测试之外使用。你的用例是什么?数据应该由执行者并行读取、处理和写入。如果需要,您可以在其他地方使用生成的数据。 -
用例是将其写入数据库。
-
为什么 '.collect()' 不能用于实验和测试?超越它的哲学是什么?有没有关于其最佳实践的文章/博客?
-
Spark 可以并行写入数据库(Hive、JDBC 或其他)。没有必要收集。 Spark 最适合用于大规模并行。使用 collect 时,您会将所有数据拉到一个节点中,从而杀死所有并行性,而简单的 Python 脚本很容易超越。
标签: pyspark rdd batch-processing