Spark分区与Task关系详解:并行度与分区数调优指南

发布时间:2026/10/6 3:40:36
Spark分区与Task关系详解:并行度与分区数调优指南 先别急着搜源码先说结论这个问题比你想的复杂一点点但也就复杂那么一点点。我在答疑群和后台经常看到有人问“我把spark.default.parallelism设成18了为什么UI上task数还是50多个”“我的executor有8个core为什么task跑不满”这类问题归根结底就一句话你没分清partition和task到底谁决定谁也没搞明白它们在各阶段分别被什么参数影响。这篇文章我把这两件事彻底拆开讲。从数据进入Spark那一刻开始一路跟踪到task在executor上被调度把每个阶段的“分区数怎么来”和“并行度怎么变”捋清楚。适合正要调优的、准备面试的、还有读完源码但没串成线的同学。1. 先把概念对齐partition是数据切块task是计算单元这两个词在Spark语境里经常被混用但它们不在一个层面上。Partition是数据层面的概念。RDD被切成若干块每一块数据独立存储在某个executor的某个block里也可能在磁盘。分区是数据的物理切分方式是Spark实现分布式计算的前提——没有切分就没有并行。Task是计算层面的概念。一个分区对应一个tasktask是运行在executor里的最小计算单元。Spark把一个stage的rdd计算逻辑复制成N份每份任务处理自己对应的分区数据。打个比方partition是“一大块猪肉切成的肉片”task是“每个厨师拿到一片肉去炒”。厨师的数量取决于肉切了多少片而不是取决于厨房里有几个灶台。1.1 那个绕不开的映射关系一个分区分一个task这是整篇文章的地基先刻在脑子里一个partition在某个stage里对应一个task。Spark UI上显示的task数通常就等于当前stage中所有输入分区数的总和。也就是说task的个数在绝大多数情况下是被partition数量锁死的不是被executor数量锁死的也不是被core数量锁死的。那“并行度”这个词到底指什么在Spark官方文档里并行度通常指“同一时刻正在运行的task数量”或者说整个集群能同时处理的task数。这个上限由executor数量、每个executor的core数、每个task消耗的CPU数共同决定。所以你会遇到两种情况分区数远大于并行度上限task排队执行每个task都小但总共要跑很久。分区数远小于并行度上限一堆core闲着每个task都巨大跑得慢且容易OOM。真正舒服的状态是分区数 ≈ 并行度上限的2~4倍让每个task不至于太大同时又能让所有core尽量保持忙碌。很多人的误区就在这里以为调大executor的core数就能提高并行度。不对那只是提高了“同时能跑多少task”的上限真正的并行粒度还是由分区数决定。1.2 为什么一把梭“分区数越大并行度越高”是错的分区数多少会影响task执行代价。每个task有固定的调度开销、序列化开销、结果从executor传到driver的成本。分区数上到几千几万这些开销会被无限放大。反过来分区数太少也不行。每个task处理的数据量太大会导致GC压力大、数据溢出到磁盘、单个task执行时间过长拖慢整体。所以后面我们讨论的一切本质都是“如何让分区数匹配你的资源”。2. 分区数的决定链路从源头数据一路算到算子分区的数量是动态演变的——从头到尾不断被各种操作改变。我们逐个阶段看。2.1 读取阶段HDFS、本地文件、并行集合各玩各的规则输入源不同初始分区规则完全不同。下面这张表是基本对照。数据源分区数决定方式备注sc.textFile(hdfs://...)默认按HDFS block数量决定一个block一个partitionblock默认128MB可通过第二个参数minPartitions设置“最小分区数”sc.wholeTextFiles(hdfs://...)一个文件一个partition对小文件极不友好sc.parallelize(collection)默认按spark.default.parallelism也可通过第二个参数numSlices指定numSlices就是分区数会覆盖默认值sc.textFile(s3://...)按对象列表每个对象对应一个或多个partition取决于s3a/hadoop的split逻辑JDBC读取JdbcRDD由numPartitions参数指定默认值通常为3每个分区执行一段SQL查询最典型的是HDFS文件。Spark底层用InputFormat去切分文件TextInputFormat默认一个block对应一个InputSplit一个split对应一个partition。当文件是128MB的10个块时初始就是10个分区也就是10个task。这里有个关键细节textFile的第二个参数minPartitions名字叫“最小分区数”不是精确的分区数。如果你传1最后可能出来2个分区因为Hadoop的split计算逻辑有自己的校验。想精确控制那就别用textFile用hadoopFile传自定义split或者先repartition。另外一个高频陷阱是小文件。如果HDFS上有几千个几十KB的小文件textFile会为每个文件至少生成一个分区导致初始分区数爆炸Spark UI上几千个task在飞调度开销和元数据开销直接把这个作业拖垮。解决办法是在读取后立刻coalesce或者更优雅一点在上游把文件先合并成大文件。2.2 RDD算子阶段窄依赖保持分区宽依赖重置分区读完数据之后分区数还要看你做了什么操作。这里要把“窄依赖”和“宽依赖”搞清楚。窄依赖算子map、filter、flatMap、mapPartitions等每个父分区只被子分区使用分区数保持不变。所以就算你写一百个map分区数还是跟输入一样。很多人混淆的点在于我用了map为什么task数没变因为窄依赖根本不会触发shuffle分区数自然不变。宽依赖算子reduceByKey、groupByKey、join、distinct、repartition等父分区的数据要被重新分配到不同的子分区会产生shuffle分区数会被重置。reduceByKey(func, numPartitions)这种带numPartitions参数的算子你传多少最终分区数就是多少。不传的话走默认逻辑——见下一节。join也是一样的套路可以传numPartitions。不传就取两个父RDD中最大的分区数。partitionBy(new HashPartitioner(n))直接把RDD按Key的Hash分区为n份最终分区数固定为n。2.3 三个常被混淆的参数default.parallelism、shuffle.partitions、手动指定这是最容易搞混的三个东西我把它们并排放一起。spark.default.parallelism作用于RDD层面。它决定了parallelize这类并行集合的默认切片数以及reduceByKey等shuffle算子在不显式指定分区数时的新分区数。它的默认值local模式下本地CPU核数YARN/Standalone模式下所有executor的总cores数且最小为2有个细节要知道YARN模式下这个默认值有时并不合理。比如你有10个executor、每个4核它默认也就40。如果你的数据量特别大40个分区在大量join之后会显得力不从心所以生产环境最好显式设置。spark.sql.shuffle.partitions作用于Spark SQL/DataFrame层面。SQL/Hive表的join、group by、distinct这些操作shuffle后的分区数由它控制默认值是200。为啥是200因为Spark早期认为200是个安全的中间值既能保证一定的并行度又不会产生过多小task。但这个值在现代硬件条件下偏保守了尤其是两个大表join的场景200个分区可能导致每个分区处理几个GB数据严重的OOM和磁盘溢写。大集群建议按数据量调整到500~2000。手动指定rdd.reduceByKey(__, 100)、df.repartition(500)、df.coalesce(50)、rdd.partitionBy(new HashPartitioner(100))。这种指定是最高优先级的直接覆盖上面两个默认参数。注意spark.sql.shuffle.partitions管不到RDD API。你要是在RDD上用了reduceByKey又不传参数走的是spark.default.parallelism。这个坑我见太多了一个项目里既写RDD又写SQL调了半天spark.sql.shuffle.partitions没反应最后发现RDD那边根本没吃这个参数。2.4 shuffle是个重头戏宽依赖如何重置分区数shuffle的一次完整流程可以简化成三句话上游map侧把数据按目标分区号写入本地磁盘文件下游reduce侧拉取对应分区的数据reduce侧的分区数决定了下游stage的task数。那“目标分区号”怎么算Key的Hash值对分区数取模。Hash分区是默认算法存在数据倾斜的风险——某个Key的Hash值大量落在同一分区。这也是为什么很多人在repartition之后看到某个task数据量比其他task大几十倍那就是Hash分区的某种意义上的“公平但不均”。另一种分区器是RangePartitioner用于sortByKey。它会按照Key的顺序范围去划分分区保证每个分区的数据量大致均匀。但代价是需要采样统计Key分布构造分区器本身要额外消耗一个job的算力。如果shuffle算子显式指定了分区数shuffle后的RDD分区数就是你指定的值。没指定RDD默认走spark.default.parallelismSQL默认走spark.sql.shuffle.partitions。在这里我想提醒一个常见的性能事故小表join大表结果分区数直接继承了大表的分区数而小表那边每个分区会被广播或者拉取很多次。遇到这种情况给小表先repartition(一个合理值)能显著降低shuffle数据量。3. 并行度task数量的真实决定链路从DAG到Executor很多人查task数量只知道看Spark UI总数不知道这个数字是怎么一层层被决定的。捋清楚后面这些你就不会再犯“改了executor配置但task数不涨”这种错误。3.1 惊掉下巴的事实task数通常是当前stage的所有分区数相加一个Spark作业会切分成多个stage每个stage底层的RDD分区数就是该stage的task数。闭包、计算的具体过程是这样每个RDD有一个partitions数组记录全部分区。每个stage的rdd.partitions.length就是task数的基础。ShuffleMapStage每个partition在map侧产生一个ShuffleMapTask。ResultStage每个partition产生一个ResultTask也就是最后输出给driver的task。举个实际例子你读了一个100分区的HDFS文件做map操作分区数不变还是100再reduceByKey重新分区成50。整个过程两个stageStage1map侧ShuffleMapStage100个taskStage2reduce侧ResultStage50个taskSpark UI上你会看到job总task数千变万化就是因为不同stage的task数不一样。如果你看到UI上某个stage是100个task下一个stage变成50个不要以为是配置生效了那是shuffle重置了分区数。有一种特殊情况同一stage可能有部分task失败重试UI上task数会临时增加。另外如果开了spark.speculation推测执行慢task会被复制一份去别的executor跑UI上也会多出task。这个多出来的task不会增加数据切片纯粹是“备份跑”的性质。3.2 调度层面cores、executors、CPU资源如何限制“同时并行”的task数分区数决定了Spark要生成多少个task但同一时刻能并行跑几个由资源决定。单个executor可以并行执行的task数每个executor的可用core数 ÷ 每个task占用的core数spark.task.cpus默认是1也就是说默认一个core同时跑一个task。如果一个executor有4个core它同一时刻并行跑4个task。如果spark.task.cpus2那一个executor只能同时跑2个task相当于每个task要独占2个核适合CPU密集或者需要避免频繁上下文切换的场景。整个集群的并行度上限executor总数 × 每个executor可用core数 ÷ spark.task.cpus这就是调度器控制并发队列长度的依据。FIFO调度器下所有已提交的task进入待执行队列只要有空闲core就拉一个去执行FAIR调度器则按pool分配资源比例。到这里你应该明白一件事task总数分区数和“同时跑的task数”是两个独立的维度。前者取决于数据和算子后者取决于资源。你永远不可能靠增加executor来增加task数你增加的executor只会让你同一时刻跑更多的task。如果分区数只有20给你100个core同一时刻最多也只能跑20个task剩下80个core闲着。3.3 Spark UI里的Task表到底在告诉你什么打开Spark UI进入某个stage点“Tasks”标签页你会看到一张表格里面有Task Index、Executor ID、Duration、GC Time、Shuffle Read、Shuffle Write等列。这里最有价值的是三列Task Index分区的序号从0开始。如果某个index的Duration特别长对应的partition数据量多半很大——这就是数据倾斜的直观证据。Shuffle Read/Shuffle Write读和写的数据量。同一个stage里不同task的Shuffle Write差异巨大说明上游分区数据不均。Locality LevelPROCESS_LOCAL是最好状态数据就在同一executor内存里NODE_LOCAL次之RACK_LOCAL要跨机架跨网络读数据自然慢。我排查性能问题的一个固定动作直接看Tasks表按Duration排序看中位数和最大值的差距。如果中位数是2秒、最大是5分钟不用再猜了分区数据分布出了问题优先处理倾斜而不是盲目加资源。4. 实战中的常见误区参数设了为什么没生效这部分记录我实际遇到过的问题比理论本身更值得抄走。4.1 误区一把spark.default.parallelism当成万能并行度开关之前说过这个参数管的是RDD API里shuffle后的默认分区数、parallelize默认切片数。但它管不到已经存在的RDD窄依赖不会改变分区数Spark SQL的shuffle分区数走spark.sql.shuffle.partitions显式指定了分区数的算子举个我遇到的真实案例有个同事在conf里设了spark.default.parallelism1000然后跑了SQL join死活看到shuffle后只有200个分区。他以为是集群大小问题折腾了半天最后发现SQL走的是另一个参数。改成spark.sql.shuffle.partitions1000立刻生效。4.2 误区二coalesce(n)和repartition(n)可以混用repartition(n)是full shuffle会把整个RDD的数据全部重新打乱再分n份。适合增大分区数、或者让你重新分布数据到尽量均匀。coalesce(n, shufflefalse)不会打乱数据只是把多个分区的数据在executor本地合并到一起。适合在stage末尾裁剪分区数减少输出文件数。问题在于coalesce的合并逻辑是本地的不能跨executor移动数据。如果n小于当前executor数量会导致某些executor完全没有数据分区另外一些executor反而堆很多分区。举个例子你有10个executor当前800个分区coalesce(5)之后极大概率是分给了前5个executor因为同一executor上有多个分区时优先本地合并后5个executor空转。如果你想减少分区数且保证数据尽量均衡正确姿势是coalesce加shuffletrue或者直接用低配版repartition。当然代价是多一点shuffle开销但换来的是均衡。另外一个坑coalesce(n)传的n大于当前分区数时不会增加分区最多保持原样。想增分区只能repartition。4.3 误区三以为task数越小越好或越大越好这个我要说得直接一点都不是。task数过大调度开销大、每个task的计算时间太短几个毫秒大量时间被浪费在task启动/结果回收上。我之前压测过一个纯map作业本来是1000个分区跑了30秒改成10000个分区反而跑了55秒——纯粹是调度和序列化开销翻了几倍。task数过小每个task处理数据量太大GC压力剧增可能频繁触发磁盘溢写整个stage的执行时间被少数几个task卡住。我见过一个join作业只有5个分区每个分区3GB最后2个executor OOM作业失败。repartition(20)之后5分钟跑完。我常用的经验区间每个task处理的数据量控制在100MB到1GB之间再结合executor并发度去定分区数。这个区间不是绝对的但作为起点很好用。4.4 一个完整的排查案例为什么task只有2个但executor有8个之前有个同学发了个截图16个executor每个8个core跑一个读取400MB文本的作业UI上task只有2个且所有executor都在等这两个task跑完。看一眼代码val rdd sc.textFile(/data/access.log, 2)他手动指定了minPartitions2以为这只是“最少”的意思实际Spark按这个值生成了2个分区每个200MB。加上后续窄依赖全保持2个分区整个作业从头到尾就是2个task。至于task为什么只落在2个executor上而不是所有16个执行为什么没跑满所有core因为task太少2个task只需要2个core就能同时跑剩下的core自然闲着。修复val rdd sc.textFile(/data/access.log) // 走block默认400MB大概拆分4个分区 // 或者 val rdd sc.textFile(/data/access.log, 32) // 直接调到目标并行度如果文件特别大建议按“executor总core数×2~4”去估算分区数不要硬算成block数。block数只适合中小数据量场景。还有一次更隐蔽的代码里没指定分区数HDFS上文件也就2个block凑巧每个block 128MB只有两个文件于是task数就是2。这个不是配置问题是源数据文件太少。处理方式要么把小文件合并成大文件后在读取时保持合理分区数要么读取后立即repartition。5. 调优实操怎么确定合理的分区数与并行度配置这一节直接给步骤和公式按下面的流程走基本不会踩大坑。5.1 先盘点资源再定分区数基准值第一步是算你集群到底能同时跑多少task。公式同时并行task数 executor数量 × 每个executor的核数 ÷ spark.task.cpus举个例子100个executor每个4核默认spark.task.cpus1同时并行400个task。然后分区数基准值我建议取同时并行task数的2到4倍。理由是比并行数小core利用率不够比并行数大太多单个task太小调度开销膨胀于是上面这个集群分区数区间就是800~1600。后续shuffle算子、SQL的shuffle.partitions都朝这个区间靠。5.2 按不同阶段去配套设置读取、shuffle、输出读取阶段HDFS文件如果文件布局健康block数接近目标分区数直接什么参数都不设按block分区就很合理。小文件很多读取后先coalesce(目标分区数)注意上面说的局限性必要时用repartition。同时读多个小文件目录用wholeTextFiles之后分区数会按文件数爆炸建议直接repartition。shuffle阶段RDD APIreduceByKey(_, _, numPartitions)里显式传值。我不建议依赖全局默认因为不同stage的数据量可能差异很大。Spark SQLjob启动前执行spark.sql(SET spark.sql.shuffle.partitions800)或者写conf文件。注意这个设置是session级的在同一个spark-shell里改了会影响后续所有SQL作业。输出阶段如果目标是“每个文件大小接近目标块大小比如128MB”用coalesce控制输出文件数。如果目标是“下游读取速度最优”输出文件数最好跟下游HDFS block对齐不要出现一个文件几十GB或者几KB的极端情况。5.3 关于数据倾斜的补充分区不均时先别急着调全局一个非常常见的场景不是分区少而是某个分区数据特别大导致其他task都跑完了就那一个task还在跑。这时候你把全局spark.sql.shuffle.partitions调大几百倍可能把倾斜分区分开了但也让所有正常分区的task小得离谱白浪费调度时间。更好的处理顺序是先定位是不是真的倾斜——看Tasks表里Duration是否一骑绝尘看Shuffle Read是否集中。如果倾斜来自某个Key优先做Key加盐Salt或两阶段聚合。如果倾斜来自数据分布本身就偏考虑改用Range分区器或者自定义Partitioner而不是单纯调大分区数。实在不行再考虑调spark.sql.shuffle.partitions用分区数的增大来摊薄倾斜。另外spark.sql.adaptive.enabled在Spark 3.x默认开AQE自适应查询执行会在shuffle之后检测实际分区大小动态合并小分区、分割倾斜分区。开了AQE之后spark.sql.adaptive.coalescePartitions.enabled也会帮助你自动减少无用的小task。这个功能建议别关能让不少手动调优变成自动调优。5.4 给小白的一个“不会错”的起步配置如果你的作业没那么复杂想先跑起来再优化可以用这套起步值每个executor 4核、8GB内存executor数量按集群实际节点数配置spark.default.parallelism executor总核数 × 2spark.sql.shuffle.partitions executor总核数 × 2所有RDD API的shuffle算子显式传同样的numPartitionsSQL用AQE默认配置即可这套配置不算最优但基本不会出现“几十个core空转”或者“task太大OOM”这种低级问题。跑起来之后再用UI观察微调分区数。我实际做压测时往往会准备三组配置分别按并行度的1倍、2倍、4倍设置分区数用同一份数据跑三遍对比耗时。多数情况下2倍左右综合最优但数据倾斜严重的作业可能会在4倍时表现更好——这就是为什么要实测而不是纯靠理论拍脑袋。