【问题标题】:Spark Creates Less Partitions Then minPartitions Argument on WholeTextFilesSpark 在 WholeTextFiles 上创建更少的分区然后 minPartitions 参数
【发布时间】:2019-01-08 18:21:10
【问题描述】:

我有一个文件夹,里面有 14 个文件。我在一个集群上使用 10 个执行器运行 spark-submit,它的资源管理器为 yarn。

我这样创建我的第一个 RDD:

JavaPairRDD<String,String> files = sc.wholeTextFiles(folderPath.toString(), 10);

但是,files.getNumPartitions()随机给我 7 或 8 个。然后我不在任何地方使用合并/重新分区,我用 7-8 个分区完成了我的 DAG。

据我所知,我们给出的参数是“最小”分区数,那么为什么 Spark 将我的 RDD 划分为 7-8 个分区?

我也用 20 个分区运行相同的程序,它给了我 11 个分区。

我在这里看到了一个话题,但它是关于“更多”分区的,这对我没有任何帮助。

注意:在程序中,我读取了另一个包含 10 个文件的文件夹,Spark 成功创建了 10 个分区。在这个成功的工作完成后,我运行上述有问题的转换。

文件大小: 1)25.07 KB 2)46.61 KB 3)126.34 KB 4)158.15 KB 5)169.21 KB 6)16.03 KB 7)67.41 KB 8)60.84 KB 9)70.83 KB 10)87.94 KB 11)99.29 KB 12)120.58 KB 13)170.43 KB 14)183.87 KB

文件在 HDFS 上,块大小为 128MB,复制因子 3。

【问题讨论】:

    标签: apache-spark hdfs hadoop-yarn distributed-computing partitioning


    【解决方案1】:

    如果我们有每个文件的大小会更清楚。但是代码不会出错。我根据 spark 代码库添加这个答案

    • 首先,ma​​xSplitSize的计算取决于目录大小最小分区 传入wholeTextFiles

          def setMinPartitions(context: JobContext, minPartitions: Int) {
            val files = listStatus(context).asScala
            val totalLen = files.map(file => if (file.isDirectory) 0L else file.getLen).sum
            val maxSplitSize = Math.ceil(totalLen * 1.0 /
              (if (minPartitions == 0) 1 else minPartitions)).toLong
            super.setMaxSplitSize(maxSplitSize)
          }
          // file: WholeTextFileInputFormat.scala
      

      link

    • 根据maxSplitSize splits(Spark 中的分区)将从源中提取。

          inputFormat.setMinPartitions(jobContext, minPartitions)
          val rawSplits = inputFormat.getSplits(jobContext).toArray // Here number of splits will be decides
          val result = new Array[Partition](rawSplits.size)
          for (i <- 0 until rawSplits.size) {
            result(i) = new NewHadoopPartition(id, i, rawSplits(i).asInstanceOf[InputSplit with Writable])
          }
          // file: WholeTextFileRDD.scala
      

      link

    CombineFileInputFormat#getSplits 课堂上提供有关读取文件和准备拆分的更多信息。

    注意:

    我在这里将 Spark 分区称为 MapReduce 拆分,称为 Spark 从 MapReduce 借用输入和输出格式化程序

    【讨论】:

    • 我正在提供问题中文件的大小。
    猜你喜欢
    • 1970-01-01
    • 2020-09-29
    • 2018-04-18
    • 1970-01-01
    • 2020-07-14
    • 2019-08-04
    • 1970-01-01
    • 2016-02-27
    • 2022-10-05
    相关资源
    最近更新 更多