【问题标题】: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() 以增加/减少分区。

      【讨论】:

        猜你喜欢
        • 2020-01-01
        • 2022-12-17
        • 2020-04-02
        • 1970-01-01
        • 1970-01-01
        • 2016-04-23
        • 1970-01-01
        • 2018-01-10
        • 2016-02-22
        相关资源
        最近更新 更多