【问题标题】:load parquet file and keep same number hdfs partitions加载 parquet 文件并保持相同数量的 hdfs 分区
【发布时间】:2019-10-29 07:41:08
【问题描述】:

我有一个镶木地板文件/df 保存在具有 120 个分区的 hdfs 中。 hdfs上每个分区的大小在43.5M左右。

总大小

hdfs dfs -du -s -h /df
5.1 G  15.3 G  /df
hdfs dfs -du -h /df
43.6 M  130.7 M  /df/pid=0
43.5 M  130.5 M  /df/pid=1
...
43.6 M  130.9 M  /df/pid=119

我想将该文件加载到 Spark 中并保持相同数量的分区。 但是,Spark 会自动将文件加载到 60 个分区中。

df = spark.read.parquet('df')
df.rdd.getNumPartitions()
60

HDFS 设置:

'parquet.block.size' 未设置。

sc._jsc.hadoopConfiguration().get('parquet.block.size')

什么都不返回。

'dfs.blocksize' 设置为 128。

float(sc._jsc.hadoopConfiguration().get("dfs.blocksize"))/2**20

返回

128

将这些值中的任何一个更改为更低的值不会导致 parquet 文件加载到与 hdfs 中相同数量的分区中。

例如:

sc._jsc.hadoopConfiguration().setInt("parquet.block.size", 64*2**20)
sc._jsc.hadoopConfiguration().setInt("dfs.blocksize", 64*2**20)

我意识到 43.5 M 远低于 128 M。但是,对于这个应用程序,我将立即完成许多转换,这将导致 120 个分区中的每一个都更接近 128 M。

我正在努力避免在加载后立即在应用程序中重新分区。

有没有办法强制 Spark 加载与存储在 hdfs 上的分区数相同的 parquet 文件?

【问题讨论】:

  • 如果您尝试设置此参数会怎样?到 43.5MB(43500000 字节)spark.conf.set("spark.files.maxPartitionBytes", 43500000)
  • 不。仍然将其拉入 60 个分区。

标签: apache-spark hadoop pyspark apache-spark-sql parquet


【解决方案1】:

首先,我将首先检查 Spark 如何将数据拆分为分区。 默认情况下,它取决于数据和集群的性质和大小。 这篇文章应该会为您提供为什么您的数据框被加载到 60 个分区的答案:

https://umbertogriffo.gitbooks.io/apache-spark-best-practices-and-tuning/content/sparksqlshufflepartitions_draft.html

一般来说 - 它的 Catalyst 负责所有优化(包括分区数量),所以除非真的有充分的理由进行自定义设置,否则我会让它完成它的工作。如果您使用的任何转换很宽,Spark 无论如何都会打乱数据。

【讨论】:

    【解决方案2】:

    我可以使用spark.sql.files.maxPartitionBytes 属性在导入时将分区大小保持在我想要的位置。

    spark.sql.files.maxPartitionBytes 属性的 Other Configuration Options documentation 声明:

    读取文件时打包到单个分区的最大字节数。此配置仅在使用 Parquet、JSON 和 ORC 等基于文件的源时有效。

    示例(其中spark 是有效的SparkSession):

    spark.conf.set("spark.sql.files.maxPartitionBytes", 67108864) ## 64Mbi
    

    为了控制转换过程中的分区数量,我可以设置spark.sql.shuffle.partitionsdocumentation 表示:

    配置在为联接或聚合打乱数据时要使用的分区数。

    示例(其中spark 是有效的SparkSession):

    spark.conf.set("spark.sql.shuffle.partitions", 500)
    

    另外,我可以设置spark.default.parallelismExecution Behavior documentation 声明:

    当用户未设置时,由 join、reduceByKey 和并行化等转换返回的 RDD 中的默认分区数。

    示例(其中spark 是有效的SparkSession):

    spark.conf.set("spark.default.parallelism", 500)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-12-02
      • 2023-03-20
      • 1970-01-01
      • 2019-07-29
      • 2017-11-08
      相关资源
      最近更新 更多