【问题标题】:How to collect elements in a RDD by batches如何批量收集RDD中的元素
【发布时间】:2021-10-12 10:25:46
【问题描述】:

我有一个 pyspark RDD,它有大约 200 万个元素。我不能一次收集它们,因为它会导致OutOfMemoryError 异常。

如何分批收集?

这是一个潜在的解决方案,但我怀疑有更好的解决方案:收集一个批次(使用 takehttps://spark.apache.org/docs/3.1.2/api/python/reference/api/pyspark.RDD.take.html#pyspark.RDD.take),然后从该批次中的 RDD 中删除所有元素(使用 filterhttps://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


【解决方案1】:

我不确定它是否是一个好的解决方案,但您可以使用索引压缩您的 rdd,然后对该索引进行过滤以批量收集项目:

big_rdd = spark.sparkContext.parallelize([str(i) for i in range(0, 100)])
big_rdd_with_index = big_rdd.zipWithIndex()
batch_size = 10
batches = []
for i in range(0, 100, batch_size):
  batches.append(big_rdd_with_index.filter(lambda element: i <= element[1] < i + batch_size).map(lambda element: element[0]).collect())
for l in batches:
  print(l)

输出:

['0', '1', '2', '3', '4', '5', '6', '7', '8', '9']
['10', '11', '12', '13', '14', '15', '16', '17', '18', '19']
['20', '21', '22', '23', '24', '25', '26', '27', '28', '29']
['30', '31', '32', '33', '34', '35', '36', '37', '38', '39']
['40', '41', '42', '43', '44', '45', '46', '47', '48', '49']
['50', '51', '52', '53', '54', '55', '56', '57', '58', '59']
['60', '61', '62', '63', '64', '65', '66', '67', '68', '69']
['70', '71', '72', '73', '74', '75', '76', '77', '78', '79']
['80', '81', '82', '83', '84', '85', '86', '87', '88', '89']
['90', '91', '92', '93', '94', '95', '96', '97', '98', '99']

【讨论】:

  • 有趣,谢谢!
  • 它需要评估性能和内存,但我喜欢这个想法
猜你喜欢
  • 1970-01-01
  • 2023-02-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-01-29
  • 1970-01-01
相关资源
最近更新 更多