【问题标题】:Spark coalescing on the number of objects in each partitionSpark合并每个分区中的对象数量
【发布时间】:2018-12-18 00:19:25
【问题描述】:

我们开始在我们的团队中试验 spark。 在 Spark 中完成 reduce 工作后,我们希望将结果写入 S3,但是我们希望避免收集 spark 结果。 目前,我们正在将文件写入 RDD 的 Spark forEachPartition,但这会导致很多小文件。我们希望能够将数据聚合到几个文件中,这些文件按写入文件的对象数量进行分区。 例如,我们的总数据是 1M 个对象(这是恒定的),我们想生成 400K 个对象文件,而我们当前的分区产生大约 20k 个对象文件(每个作业变化很大)。理想情况下,我们希望生成 3 个文件,每个文件包含 400k、400k 和 200k,而不是 20K 对象的 50 个文件

有人有好的建议吗?

我的想法是让每个分区处理它应该写入哪个索引,假设每个分区将大致产生相同数量的对象。 例如,分区 0 将写入第一个文件,而分区 21 将写入第二个文件,因为它假定对象的起始索引是 20000 * 21 = 42000,它大于文件大小。 分区 41 将写入第三个文件,因为它大于 2 * 文件大小限制。 不过,这并不总是会导致完美的 400k 文件大小限制,更多的是近似值。

我知道有合并,但据我了解,合并是根据想要的分区数减少分区数。我想要的是根据每个分区中的对象数量合并数据,有没有好的方法呢?

【问题讨论】:

  • 为什么不将其存储到镶木地板或 ORC 中?您可以使用.repartition 而不是.coalesce 来确定您想要的文件的确切数量。
  • 您能否详细说明将其存储到 ORC 或 parquet 中?
  • 另外,我可以确定文件的确切数量,但是,我并不关心生成了多少文件。我更关心生成的文件有多大

标签: apache-spark


【解决方案1】:

您要做的是将文件重新分区为三个分区;数据将被拆分为每个分区大约 333k ​​条记录。分区将是近似的,每个分区不会完全是 333,333。我不知道如何获得您想要的 400k/400k/200k 分区。

如果你有一个 DataFrame `df',你可以重新分区成 n 个分区

df.repartition(n)

由于您需要每个分区的最大数量或记录,我建议您这样做(您不指定 Scala 或 pyspark,所以我将使用 Scala;您可以在 pyspark 中执行相同操作):

val maxRecordsPerPartition = ???
val numPartitions = (df.count() / maxRecordsPerPartition).toInt + 1
df
    .repartition(numPartitions)
    .write
    .format('json')
    .save('/path/file_name.json')

这将保证您的分区小于 maxRecordsPerPartition。

【讨论】:

  • 是的,这就是我们最终做的事情
【解决方案2】:

我们决定只考虑生成的文件数量,并确保每个文件包含的行项目少于 100 万个

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-19
    • 1970-01-01
    • 2015-01-01
    • 2016-10-02
    • 2017-09-11
    • 2017-09-29
    相关资源
    最近更新 更多