☰
Spark分区数与并行度详解:从Partition到Task的调度逻辑与调优实践
2026/10/5 13:41:39 网站建设 项目流程

写 Spark 调优的人经常会碰到这么一个问题:我明明在代码里写了repartition(1000),为什么 Spark UI 里的 task 数却不是 1000?或者反过来,我以为并行度由 Executor 核数决定,为什么某个 stage 里的 task 数又会超过核数?这篇内容我会把 Spark 的 partition(分区数)、并行度(通常表现为 RDD 中 task 的个数)之间的关系彻底理一遍,包括它们分别由哪些因素决定,以及实战里怎么调才不浪费资源。适合刚接触 Spark 的开发者,也适合那些已经写过很多 Spark SQL、但对 Spark 底层任务调度仍有疑惑的人。

1. 从 RDD partition 到 task:先理解 Spark 的并行执行逻辑

1.1 Partition 是数据分片,不是“线程”

Spark 中一个 RDD 实际上就是一组 partition 的集合。如果用一句话说,partition 就是“分布在集群不同 Executor 上的数据分片”。每个 partition 在物理上可能是一个文件片段、一段内存里的对象数组、或者经过 shuffle 后落到某个 Executor 上的记录集合。RDD 的抽象让 Spark 把大数据集合当成逻辑上的单个集合处理,但真正计算时,数据是被拆散到各个 partition 里的。

很多人容易把 partition 和线程、进程混在一起,这是后续理解 task 的最大障碍。partition 只是“数据的单元”,不是“执行的单元”。真正执行运算的单位是 task,一个 task 负责处理一个 partition 上的数据。所以一个 stage 中如果有 100 个 partition,就一定会生成至少 100 个 task(map 阶段)。分区数据量大、数据分布不均,最终影响的是每个 task 的处理耗时,而 task 的数量和分区数总是静态绑定的。

1.2 并行度、Task 与 CPU 核数的三角关系

“并行度”这个词在 Spark 里有多个层面。宏观上,并行度代表整个应用同时执行的 task 数量上限,这受 Executor 总核数制约;微观上,讨论某个 stage 时,并行度往往指该 stage 的 task 总数,这个数由分区数决定。两者的关系可以概括为“分区分出任务,核数决定并发”。

可以用厨房的例子帮助记忆:partition 是切好的菜品,task 是每个菜品从切配到出锅的执行单,Executor 里的 CPU 核就是灶头。切多少份菜,大概会产生多少张执行单;但灶头只有那么多,同一时刻只能同时炒若干道菜,剩下的菜只能排队。理解了这一点,就不会再纠结“明明分区数很多,为什么并行程度看起来不高”了。

概念本质主要决定因素
Partition数据分片的大小和个数数据源大小、创建方式、repartition/coalesce、shuffle 配置
Task一个 partition 上的一次计算执行单元当前 stage 的分区数
并行度(并发 task 上限)同一时刻能跑的 task 数Executor 核数、动态分配、调度资源

1.3 Stage 中 task 数量的来源

Spark 的 job 会被 DAG 调度器按宽依赖切分成多个 stage。stage 之间通过 shuffle 连接,stage 内部是连续的窄依赖计算。因此每个 stage 的 task 数并不一定全局相同:map 端 stage 的 task 数等于上游 RDD 的分区数;reduce 端 stage 的 task 数等于 shuffle 后生成的新分区数。

举一个最简单的例子:

rdd1 = sc.parallelize(range(1000), 100) rdd2 = rdd1.map(lambda x: (x % 10, x)) rdd3 = rdd2.groupByKey(50) print(rdd3.count())

rdd2所在 stage(map 阶段)有 100 个分区,所以 Spark UI 会显示 100 个 task;groupByKey(50)会对数据做哈希重分区,产生 50 个 shuffle 输出分区,因此后面的 reduce 阶段 task 数是 50。这就是“task 个数等于 stage 对应 RDD 分区数”的直接证据。

2. 分区数由什么决定:从 RDD 到 DataFrame 的几条路径

2.1 RDD 创建时的分区数:parallelize、textFile 与 defaultParallelism

RDD 的分区数首先由创建方式决定。sc.parallelize(data, numSlices)第二个参数就是希望切成的分区数,不传时使用 SparkContext 的defaultParallelism。这个值在 local 模式下通常等于本机 CPU 核数,在集群模式下往往是集群可用总核数,你也可以通过spark.default.parallelism显式覆盖。本地跑测试时,分区数常常只有几个,你会觉得“没跑满”,这是正常的,因为本机的核数就那么多。

从文件系统创建 RDD 时则不完全一样。sc.textFile(path, minPartitions)的minPartitions只是一个下界,最终分区数由 HadoopFileInputFormat决定:HDFS 上的一个 block 通常对应一个输入分片,所以 128MB block 大小下,1GB 文件会产生 8 个左右的分区。如果你设置minPartitions=100,Spark 会尽量把文件拆成不少于 100 个分片,这意味着单个分区可能小于 block,读取会更分散。对于 gzip 这类不可分割的压缩文件,一个文件可能只能成为一个分区,这也是大压缩文件往往需要预先处理的原因。

2.2 从文件读入 DataFrame 时,分区数是怎么算出来的

现在多数数据处理走的是 DataFrame API。DataFrame 底层仍是 RDD,但分区策略交给了 Spark SQL 的文件扫描层。读取 Parquet、ORC、JSON 等文件时,总分区数大致由spark.sql.files.maxPartitionBytes(默认 128MB)控制。Spark 会计算所有文件的大小,并加上一个spark.sql.files.openCostInBytes(默认 4MB)的“打开开销”,把过小文件合并到同一个分区,从而避免生成过多微小分区。

举个例子:现在有 1000 个小文件,每个 1MB。如果不做任何合并,可能生成上千个分区,每个分区还要承担文件打开开销,调度和序列化成本很高。由于openCostInBytes的存在,Spark 倾向于把这些小文件放进更大的分区里。如果你的数据源确实有数万个小文件,可以显式调大spark.sql.files.maxPartitionBytes,比如设成 268435456(256MB)。这一层逻辑和 RDD API 的textFile不一样,很多人会忽略,而它是生产环境“分区数异常”的高频来源。

2.3 宽依赖与 shuffle 如何改写分区数

各种转换算子能直接改变分区数。repartition(n)一定触发 shuffle 并得到 n 个分区;coalesce(n)默认不触发 shuffle,但通常只能减少分区。map、filter这类窄依赖算子不会改变分区数,而groupByKey、reduceByKey、join这类宽依赖算子在 shuffle 后必然产生新分区。

这些算子中,有的允许直接传numPartitions参数,比如reduceByKey(func, numPartitions)、join(otherRDD, numPartitions)。如果你不传,RDD 层面的默认分区数是“父 RDD 最大分区数”,部分算子则使用defaultParallelism。实践里我建议凡是宽依赖都显式写分区数,否则级联 shuffle 时很容易出现“上个 stage 100 个分区,下个 stage 又变成几千个分区”的失控状态。

2.4 spark.sql.shuffle.partitions 默认 200 是把双刃剑

在 DataFrame/SQL 场景里,shuffle 后输出多少个分区由spark.sql.shuffle.partitions控制,默认 200。这个参数不会影响 RDD 算子,只管 Spark SQL 的 join、groupBy、distinct 等 shuffle 操作。默认值 200 从 Spark 早期沿用至今,因为很多作业都是几十 GB 到几百 GB 级别,200 个分区每个处理大约百 MB,刚好合适。

但“默认值适合所有作业”显然不成立。1TB 大表 join 默认只分 200 区,每个 task 要处理 5GB 数据,又慢又容易 OOM;而一个几 MB 的小 DataFrame 做 groupBy,也会生成 200 个空任务,白白浪费调度开销。我遇到生产环境都会先看数据量再手动调整,一般目标是把单个 shuffle 分区控制在 100MB~500MB 之间。在 Spark 3 中如果开了 AQE,这个参数可以不用盯那么紧,因为 AQE 会在运行时自动合并小的 shuffle 分区。

3. 并行度(task 的个数)最终取决于什么

3.1 资源核数决定“同时运行”的 task 数量上限

一个 task 最终要放到某个 Executor 的 CPU 上执行,因此任务并行度的物理上限就是所有 Executor 的可用 CPU 核数之和。公式很直接:

同时运行的最大 task 数 ≈ spark.executor.instances × spark.executor.cores

比如申请了 20 个 Executor、每个 4 核,那同一时刻最多运行 80 个 task。假设某个 stage 有 200 个分区,就会分成三波:前 80 个 task 同时跑,然后 80 个,最后 40 个。注意这里没有把超线程、CPU 共享算进去,Spark 按虚拟核心去做资源调度,所以日常监控看到的并发就是这个量级。

很多人误以为 task 数 = 分区数,所以并行度由分区数决定。严格说不对:分区数决定的是当前 stage 总共需要执行多少个 task,而资源核数决定同一时刻能并发执行多少个 task。两者都会影响整体耗时。极端情况下,哪怕你把分区数调到一万,如果 Executor 只有 8 核,同一时刻仍然只有 8 个 task 在跑,只是任务被切得更碎、调度开销更高。

3.2 并行度与分区数不一致的几种现场

我们经常在 Spark UI 里看到一种情况:stage 的 Duration 虽然短,但 task 数只有两三个,Executor 却有一堆空闲。这通常发生在读取了一个小文件,或者使用coalesce(1)写出单文件的时候。分区数小于核数,意味着大量槽位空转,资源利用率低。

另一种常见情况是分区数远大于核数,比如 2000 个分区、100 核。这样每个核要执行 20 个 task,看上去线程切换频繁了一点,但只要每个 task 数据量合理,通常也能接受。真正要避免的是单个 task 数据量超大、另一个 task 数据量几乎为零,也就是数据倾斜。这种情况下,就算分区数和核数匹配得再好,整体时间也会被最长的那几个 task 拖住。后面我会单独讲怎么排查。

3.3 动态分配与 AQE 对并行度的实际影响

共享集群里,任务状态不是一成不变的。Spark 开启spark.dynamicAllocation.enabled=true后,executor 会根据当前 job 的任务积压情况动态申请或释放。如果资源充足,并行度上限会变化。推荐设置spark.dynamicAllocation.minExecutors和spark.dynamicAllocation.maxExecutors,避免空跑时释放太快,也避免高峰时期申请过多。注意动态分配默认只适用于 YARN 和 Kubernetes 的集群模式。

AQE(自适应查询执行)是另一个影响并行度的隐藏因素。从 Spark 3.2 起spark.sql.adaptive.enabled默认开启,其中spark.sql.adaptive.coalescePartitions.enabled=true会在 shuffle 结束后,把那些数据量很小的分区自动合并,减少 reduce 端 task 数。所以你现在看到某个 SQL 的 task 数少于spark.sql.shuffle.partitions的设置值,往往不是 bug,而是 AQE 在帮你做动态调整。

4. 实战调优:分区数到底应该怎么设

4.1 先估算数据量,再倒推分区数

我见过太多人一上来就直接repartition(1000),问原因就说“task 不够并行”。正确的姿势是先明确输入数据量级和集群核数。一般经验法则是让每个分区处理 128MB~256MB 的数据(按原始文件大小估算),并尽量不要让单个 task 处理超过 1GB。比如输入是 100GB 的 Parquet,那么目标分区数大约在100GB / 200MB = 500左右。如果集群有 200 核,500 个 task 分三轮跑完,不算最理想但还能接受;如果你有 500 核,那 500 个 task 正好一轮跑完。

在核心数不变的情况下,分区数也不必严格等于核心数的整数倍。因为有的 task 跑得快,有的跑得慢,多一点分区能起到自动负载均衡的作用。通常建议分区数 = 总核数 × 2~4,也就是每个核处理 2~4 个 task,除非单个 task 数据量确实太大。这个经验在大多数 ETL 任务上是稳的。

4.2 repartition 和 coalesce 的底层行为对比

经常有人把repartition和coalesce混着用,其实两者区别很大:

操作是否 shuffle能否增加分区适用场景
coalesce(n)否(默认)通常否减少分区、降低输出小文件
repartition(n)是是增加分区、重分布数据
coalesce(n, shuffle=true)是是等价于repartition(n)

用coalesce减少分区时,如果从 1000 个分区减到 2 个,大量上游分区会被合并到少数下游分区,容易导致某个大区压力很大,因为合并时只是把多个分区串到同一个 task,并没有重新打散数据。遇到这种情况,如果必须减到很小数量,用repartition(2)虽然开销高一点,但数据分布更均匀。反过来,如果要从 10 个分区增加到 200 个,只能用repartition,它会触发 shuffle 并按 HashPartitioner 重分布。

4.3 为 Shuffle 算子单独设置 numPartitions

在 RDD API 中,几乎所有宽依赖算子都支持显式指定分区数。比如:

rdd.map(lambda x: (x % 100, 1)) .reduceByKey(lambda a, b: a + b, 50)

这样reduceByKey后的数据就只有 50 个分区,避免默认使用父 RDD 分区数。用 DataFrame 时没有这种算子级别参数,你需要临时修改配置:

spark.conf.set("spark.sql.shuffle.partitions", 100) df.groupBy("city").count().show()

要注意,spark.sql.shuffle.partitions是全局配置,会影响同一个 SparkSession 里所有后续 SQL shuffle。如果多个逻辑的运算特性差异很大,最好在代码里保存旧值,执行完再恢复,或者把这部分逻辑单独放在一个 SparkSession 中处理。我踩过不少次坑:一条语句改成 2000,后续所有 join 都变慢了。

4.4 读取小文件场景的“反向调优”

大量小文件会让分区数爆炸。这时候不需要增加分区,而是要主动减小分区。读取阶段可以调大两个参数:

spark.conf.set("spark.sql.files.maxPartitionBytes", "268435456") spark.conf.set("spark.sql.files.openCostInBytes", "8388608")

openCostInBytes可以理解为“每个文件的开销补偿”,Spark 会尽量把多个小文件塞进同一个分区,适当调大有助于合并。读取后再根据下游计算需求,用coalesce降到合理范围。如果是写入一个分区字段很少的表,写之前用repartition或coalesce控制文件数量,能显著避免 HDFS 上出现大量 KB 级碎片文件。

4.5 用 Spark UI 验证分区设置是否合理

设置完成后不要只看日志。打开 Spark UI,进入 Stage 页面,重点看三点:

  • 当前 stage 的总 task 数是否和你预期的分区数一致;
  • 每个 task 的 Shuffle Read Size / Input Size 是否悬殊;
  • Task 的 Duration 是否集中在少数大 task 上。

如果某些大 task 的输入量是其他 task 的三倍以上,说明存在数据倾斜,单纯调分区数解决不了,要配合加盐、两阶段聚合或自定义分区器。如果 task 数明显比预期少,去 Executors 页面确认当前运行中的 executor 数量和核数。很多时候你以为申请了 50 个 executor,实际因为资源不足只跑到 30 个,并行度当然上不来。

5. 常见问题与排障速查

5.1 设置了 1000 个分区,但 task 只有几十个

这个问题通常有以下几个原因。第一,当前 stage 的 RDD 分区数并没有变成 1000,比如你调用过coalesce(10),或者 AQE 在 shuffle 后自动合并了分区。第二,资源不足导致 Spark 无法同时启动更多 task,但 UI 上显示的总 task 数其实不会少,只是 task 分批执行、调度延迟变大。第三,你查看的是某个子 stage 而不是整个 job,比如某个 stage 本来就只有几十个 partition,而它在更上游的 RDD 还是 1000 个 partition。

排查时先区分“总 task 数”和“同时运行的 task 数”。如果是前者少于预期,检查rdd.getNumPartitions()或者df.rdd.getNumPartitions()。如果是后者上不去,看 Executors 页面里的 vCores 总数和当前活跃的 executor 数,再检查是否开了动态分配导致 executor 被释放。

5.2 单个 Task 耗时异常长,怎么定位是否分区分布不均

在 Spark UI 的 Stage 详情页里,把 Duration 列排序,如果出现“头部几个 task 时间占掉 stage 一半时长”,基本就是数据倾斜。用下面几步定位:

  1. 进入这些慢 task 的详情,看 Input Size / Records 是否明显大于平均值;
  2. 如果是 shuffle read 倾斜,可以看慢 task 所属分区对应的 key 分布情况;
  3. 对 join 倾斜,考虑把大表关联的小表 broadcast,或者对热点 key 加随机前缀后两阶段聚合。

值得注意的是,repartition默认使用 hash partitioner,如果某个 key 占比特别大,repartition 也无法均匀分散它,所以倾斜场景要选择加盐或 range partitioning 等方案。

5.3 写文件时小文件特别多

这是分区数过大的直接后果。如果你对输出结果调用df.write.mode("overwrite").parquet(path),写文件的 task 数等于当前 DataFrame 的分区数,每个分区会写自己的文件。如果表还有分区字段,同一目标分区内还会因为多个 task 产生更多文件。

解决办法是在写出前强制控制分区:

df.repartition(100).write.mode("overwrite").parquet("hdfs://.../output")

如果只是想减少文件数、且不要求数据严格均匀,用coalesce(100)更省。注意使用动态分区写 Hive 表时,动态分区字段会被额外拆分出多个目录,最终文件数量还会叠加分区字段数量,需要提前评估。写入后如果还是小文件太多,再考虑用OPTIMIZE或合并文件程序处理。

5.4 并行度上去了,但整体也没变快

这种情况经常出现在“计算任务不是 CPU 密集、而是外部 IO 很多”或“存在单点限制”的场景。比如一个 UDF 里做了外部接口调用,并行度越高并发请求越多,外部服务本身成了瓶颈。另一个典型问题是collect()到 driver 上处理,driver 端串行耗时盖过了 Executor 端的并行效果。优化办法是让每一步尽量保持在分布式任务内完成,减少把大量数据拉到 driver;外部调用则考虑连接池、批量接口或提前导入数据。

资源分配也需要考虑 Executor 内存。每个 partition 的数据要能被对应 Executor 装下,如果分区小而多,序列化和调度开销上升;如果分区大而少,GC 压力增大。调整分区数之后,不妨对比一下 GC 时间和 shuffle 阶段耗时,而不是只看 task 数量。

5.5 共享集群里“申请了但没拿到”的坑

实际生产中经常遇到资源配置没问题但并行度低的情况。YARN/Kubernetes 集群通常有队列限额,spark.executor.instances=50只是“最多申请 50 个”,并不代表一定立刻拿到 50 个。如果队列资源不足,Spark 只能等待。此时先看yarn application -list或对应调度平台,确认应用实际获得了多少容器。必要时减少每个 executor 的 cores/memory,让单位资源更容易被调度,或者把大 Executor 拆成多个小 Executor。

如果集群长期混部多个 Spark 应用,建议设置spark.scheduler.mode=FAIR,让不同 job 的资源竞争更公平,避免某个大应用长期占满队列,影响其他作业的并行度。

最后分享一个我自己的习惯:每次写新的 Spark 作业,启动时先把默认并行度和首份数据的分区数打出来。

print("defaultParallelism:", sc.defaultParallelism) print("numPartitions:", rdd.getNumPartitions())

跑完第一个 stage 后看一眼 Spark UI,再决定要不要调repartition或coalesce。真正把“分区数”和“并行度”分开理解后,排错和调优会顺畅很多。以上内容来自我在真实集群上踩过的坑,不一定对每个业务都最优,但覆盖了绝大多数场景的判断思路。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询