【问题标题】:How to control number of files generated while setting large partitions in spark?如何控制在火花中设置大分区时生成的文件数量?
【发布时间】:2022-01-17 13:25:07
【问题描述】:
由于输入数据量大,我设置了 spark 的大 shuffle partitions (spark.sql.shuffle.partitions=1000)。但是,输出文件很小(~1GB),但它会创建很多小文件(3000 个文件,每个小于 1Mb)。如何将这些小文件合并为一个大文件?
还有一个问题,为什么输出文件的数量是shuffle partition数量的3倍?
【问题讨论】:
标签:
apache-spark
pyspark
apache-spark-sql
【解决方案1】:
根据 Spark 文档,spark.sql.shuffle.partitions 参数 Configures the number of partitions to use when shuffling data for joins or aggregations.。要控制输出文件的数量,请在写入输出之前使用repartition() 方法。所以是这样的:
df
.filter(...) // some transformations
.join(...)
.repartition(1) // move data into a single partition
.write
.format(...)
.save(...)
上面的 sn-p 会产生一个输出文件。
您不仅限于一次重新分区数据 - 您可以根据需要重新分区,但请记住,这是一项昂贵的操作:
df
.filter(...) // some transformations
.repartition(...) // repartition to improve join performance
.join(...)
.repartition(1) // move data into a single partition
.write
.format(...)
.save(...)
如果您想很好地解释 repartition 的工作原理,这里有一个很好的答案:
Spark - repartition() vs coalesce()
有关如何提高连接性能的更多信息,请参阅 Spark 文档:
https://spark.apache.org/docs/latest/sql-performance-tuning.html#join-strategy-hints-for-sql-queries
【解决方案2】:
因为你有大量的分区。您可能需要合并您的日期框架。 coalesce 将减少分区的数量。
val df_res = df.coalesce(10)
这应该会将输出文件的数量从 1000 个减少到仅 10 个。或者您可以 coalesce(1) 创建一个大文件。
Coalesce 使用现有分区并最小化混洗数据。结果可能大小不同。
输出文件的数量等于分区的数量。该属性 (spark.sql.shuffle.partitions) 用于连接或聚合数据的混洗。
您可以对数据框执行df.repartition() 以增加/减少分区。