Spark读取任务数确定规则:与输入分区、核心数的关联疑问
Spark读取数据时任务数的确定逻辑
核心规则拆解
Spark读取数据后的任务数,最终由处理阶段的RDD/DataFrame分区数决定,这个分区数不会直接照搬输入源的分区,而是Spark结合以下因素动态调整的结果:
- 输入源的初始分区数:比如你提到的14个6.5MB分区,这是存储层(如HDFS)的原始分区/块数,但Spark不会直接用这个数值——过小的分区会带来大量任务调度开销,Spark会自动合并小分区。
- 单个分区的最大字节限制:受参数
spark.sql.files.maxPartitionBytes控制(默认128MB),Spark会尽量让每个分区的大小接近这个值,避免出现过多极小的分区。 - 集群的可用并行能力:即总executor核数(executor数量 × 每个executor核数),Spark会尽量让最终任务数接近这个数值,确保集群资源被充分利用,避免核闲置。
结合你的测试案例分析
你的数据集总大小91MB,初始14个小分区(每个6.5MB)远小于默认的128MB分区上限,所以Spark必然会合并这些小分区,同时参考你的executor核配置:
- 测试1(2个executor,每个2核,总核数4):Spark最终生成5个任务(每个~18MB)。这个数值接近总核数4,既避免了14个小任务的调度浪费,又能让4个核同时启动任务,最后一个任务在任意核空闲后执行,保证资源利用率。
- 测试2(4个executor,每个2核,总核数8):Spark生成7个任务(每个~13MB)。这个数值更接近总核数8,相比测试1,任务数更多,能更好匹配8核的并行能力,同时每个任务的大小也在合理范围,不会因太小导致调度成本过高。
关键影响参数
spark.sql.files.maxPartitionBytes:单个分区的最大允许大小,默认128MB,是合并小分区的核心依据。spark.sql.files.minPartitionNum:最小分区数,默认等于集群总核数,确保至少有足够任务填满所有可用核。spark.default.parallelism:RDD操作的默认并行度,对于Spark SQL读取文件的场景,优先级低于前两个参数,但会作为并行度的参考值。
内容的提问来源于stack exchange,提问作者Pawel
相关产品推荐
相关产品推荐

