Spark on YARN中sc.binaryFiles()默认仅生成1个分区的原因咨询
关于Spark sc.binaryFiles()默认分区数异常的原因说明
现象说明
在Hortonworks 2.6.5版本YARN集群运行Spark 2.3.0时,调用sc.binaryFiles()读取HDFS路径下10个单文件大小4-5MB、总大小43MB的CSV文件,默认生成RDD分区数为1;同路径下调用sc.textFile()生成的RDD分区数为10,符合textFile的已知分区逻辑。
测试代码与运行结果如下:
import org.apache.log4j.{Level, Logger} import org.apache.spark.SparkContext object ReadTestYarn extends App { Logger.getLogger("org").setLevel(Level.ERROR) val sc = new SparkContext("yarn", "ReadTestYarn") val inputRDD1 = sc.textFile("hdfs:/user/maria_dev/readtest/input/*") val inputRDD2 = sc.binaryFiles("hdfs:/user/maria_dev/readtest/input/*") println("Num of RDD1 partitions: " + inputRDD1.getNumPartitions) println("Num of RDD2 partitions: " + inputRDD2.getNumPartitions) }
提交命令与输出:
[maria_dev@sandbox-hdp readtest]$ spark-submit --master yarn --deploy-mode client --class ReadTestYarn ReadTest.jar Num of RDD1 partitions: 10 Num of RDD2 partitions: 1
原理解释
这个结果是两个API底层分区逻辑的固有差异导致的,不存在配置异常:
sc.textFile()底层基于HadoopRDD实现,支持将单个文件按HDFS块边界拆分,默认分片大小和HDFS block大小对齐,分区数取输入split总数、用户指定最小分区数两者的较大值。测试场景下10个独立小文件各对应1个HDFS输入分片,总数10远大于默认最小分区参考值,最终生成10个分区,完全符合预期逻辑。sc.binaryFiles()底层基于BinaryFileInputFormat实现,设计上强制要求单个文件不可拆分,每个文件必须作为完整记录被读取,它的默认分区计算逻辑和HDFS块大小完全无关:- 未显式传入分区参数时,最小分区参考值取
sc.defaultMinPartitions,该值固定为spark.default.parallelism和2的较小值; - 遍历拿到路径下所有文件的总大小后,用总大小除以最小分区参考值,算出单个分区的目标承载容量;
- 按文件列表返回顺序,依次将文件划入当前分区,直到当前分区已承载的文件总大小超过目标容量,才切分新分区。
- 未显式传入分区参数时,最小分区参考值取
测试场景下分区数为1的具体原因:
提交任务时未手动配置spark.default.parallelism,也未给binaryFiles传入指定分区参数。HDP 2.6.5沙箱为单节点环境,默认提交任务时未开启动态资源分配,也未指定executor数量与核数,Spark计算得到的defaultParallelism为1,因此sc.defaultMinPartitions结果为1。binaryFiles计算得到单分区目标承载大小为43MB,刚好等于所有文件总大小,因此10个文件全被划入同一个分区。
如果需要调整分区数,直接在调用binaryFiles时第二个参数传入期望的分区值即可,比如sc.binaryFiles("hdfs:/user/maria_dev/readtest/input/*", 10)就能得到10个分区。
内容的提问来源于stack exchange,提问作者uhlik
相关产品推荐
相关产品推荐

