【发布时间】:2016-02-22 17:36:25
【问题描述】:
我有一个DataFrame,我需要根据特定的分区将其写入 S3。代码如下所示:
dataframe
.write
.mode(SaveMode.Append)
.partitionBy("year", "month", "date", "country", "predicate")
.parquet(outputPath)
partitionBy 将数据拆分为相当多的文件夹 (~400),每个文件夹只有一点点数据 (~1GB)。问题来了——因为spark.sql.shuffle.partitions的默认值是200,每个文件夹中1GB的数据被拆分成200个parquet小文件,总共写入了大约80000个parquet文件。由于多种原因,这不是最佳选择,我想避免这种情况。
我当然可以将spark.sql.shuffle.partitions 设置为更小的数字,比如 10,但据我了解,此设置还控制连接和聚合中随机播放的分区数,所以我真的不想更改此设置.
有谁知道是否有其他方法可以控制写入多少文件?
【问题讨论】:
-
您是否尝试过在
.write之前重新分区数据帧?乍一看,spark.sql.shuffle.partitions似乎只用于 shuffle 和 joins,但没有其他用途。否则,您应该在 partitionBy 中为额外的numParameter参数开一张票。 -
@MariusSoutier 嗯...我认为调用
repartitionbeforewrite会导致我原来的dataframe在被@重新分区之前被重新分区987654332@ 功能。将原始数据帧重新分区为仅 10 个分区肯定会导致 OOM 异常。但是,我刚刚开始测试它。更新完成后我会尽快回复。 -
@MariusSoutier 它有效!极好的。谢谢!你想把它作为回复发布吗 - 然后我会把它标记为已回答:-)
-
@GlennieHellesSindholt 我有类似的问题,但有点相反。我使用类似于您的代码来创建基于两列的分区。但我最终得到了 7-15 个文件,每个文件的大小在 30-70MB 之间。但我想增加每个分区中的文件数,并将每个文件的最大大小保持在 15MB 左右。可以这样做吗?
-
@swordfish 嗯,好吧 - 也许您可以使用
SizeEstimator来估计您的dataframe的大小,然后将其重新分区为将生成大约 15MB 文件的分区数?请注意,SizeEstimater将为您提供内存中大小的估计值,因此根据您将其存储在数据帧中的方式,您必须计算出内存中的大小如何映射到磁盘上的大小...
标签: apache-spark spark-dataframe