【问题标题】:Split files under partitions in sparkspark分区下的文件分割
【发布时间】:2020-02-21 10:26:16
【问题描述】:

我正在使用以下脚本编写分区输出。

    .write
    .format("csv")
    .partitionBy("date","region")
    .option("delimiter", "\t")
    .mode("overwrite")
    .save("s3://mybucket/myfolder/")

但是,这会导致每个分区下有 1 个文件。我想在每个分区下有多个类似大小的文件。我怎样才能达到同样的效果。我在火花2.2。

我尝试使用附加密钥作为重新分区的一部分,例如 df_input_table.repartition($"region",$"date",$"region")。但是,这会导致文件大小不同。

我想坚持使用 spark(而不是 Hive)。

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    重新分区非常昂贵,因为它会在网络上打乱数据。非常需要限制每个文件写入的最大记录数。它可以避免生成巨大的文件。在下一个版本中,Spark 提供了两种方法供用户设置限制。

    // Method 1: specify the limit in the option of DataFrameWriter API. 
    df.write.option("maxRecordsPerFile", 1000)
      .mode("overwrite").parquet(outputDirectory)
    // Method 2: specify the limit via setting the session-scoped SQLConf configuration. 
    spark.conf.set("spark.sql.files.maxRecordsPerFile", 1000)
    df.write.mode("overwrite").parquet(outputDirectory)
    

    示例 - 如果您的数据框有 10,000 条记录,并且您提供 maxRecordsPerFile = 1000,那么 spark 将创建 10 个具有相同行数的文件。

    【讨论】:

    • maxRecordsPerFile 自 Spark 2.2 起可用
    • 有问题,提到了spark 2.2。
    【解决方案2】:

    你无法控制 spark 输出文件的大小。

    repartition 不保证它仅根据键创建文件的大小假设您的文件包含 6 行,键为 A(5 rows) 和 B(1 row) 并且您将 repartitions 设置为 2 。它将创建 2 个文件,一个有 5 行,另一个文件只有 1 行。

    你可以试试这个解决方案How do you control the size of the output file?

    【讨论】:

      【解决方案3】:
      .orderBy("date","region")
      .repartition(10)
      .write
      .format("csv")
      .option("delimiter", "\t")
      .mode("overwrite")
      .save("s3://mybucket/myfolder/")
      

      您将获得 10 个几乎大小相似的文件。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-04-24
        • 2017-12-02
        • 2015-08-28
        • 1970-01-01
        • 1970-01-01
        • 2011-09-20
        相关资源
        最近更新 更多