ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Spark分区数与并行度详解:从Partition到Task的调度逻辑与调优实践

Spark分区数与并行度详解:从Partition到Task的调度逻辑与调优实践 写 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 个 taskmap 阶段。分区数据量大、数据分布不均最终影响的是每个 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所在 stagemap 阶段有 100 个分区所以 Spark UI 会显示 100 个 taskgroupByKey(50)会对数据做哈希重分区产生 50 个 shuffle 输出分区因此后面的 reduce 阶段 task 数是 50。这就是“task 个数等于 stage 对应 RDD 分区数”的直接证据。2. 分区数由什么决定从 RDD 到 DataFrame 的几条路径2.1 RDD 创建时的分区数parallelize、textFile 与 defaultParallelismRDD 的分区数首先由创建方式决定。sc.parallelize(data, numSlices)第二个参数就是希望切成的分区数不传时使用 SparkContext 的defaultParallelism。这个值在 local 模式下通常等于本机 CPU 核数在集群模式下往往是集群可用总核数你也可以通过spark.default.parallelism显式覆盖。本地跑测试时分区数常常只有几个你会觉得“没跑满”这是正常的因为本机的核数就那么多。从文件系统创建 RDD 时则不完全一样。sc.textFile(path, minPartitions)的minPartitions只是一个下界最终分区数由 HadoopFileInputFormat决定HDFS 上的一个 block 通常对应一个输入分片所以 128MB block 大小下1GB 文件会产生 8 个左右的分区。如果你设置minPartitions100Spark 会尽量把文件拆成不少于 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比如设成 268435456256MB。这一层逻辑和 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.enabledtrue后executor 会根据当前 job 的任务积压情况动态申请或释放。如果资源充足并行度上限会变化。推荐设置spark.dynamicAllocation.minExecutors和spark.dynamicAllocation.maxExecutors避免空跑时释放太快也避免高峰时期申请过多。注意动态分配默认只适用于 YARN 和 Kubernetes 的集群模式。AQE自适应查询执行是另一个影响并行度的隐藏因素。从 Spark 3.2 起spark.sql.adaptive.enabled默认开启其中spark.sql.adaptive.coalescePartitions.enabledtrue会在 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, shuffletrue)是是等价于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 一半时长”基本就是数据倾斜。用下面几步定位进入这些慢 task 的详情看 Input Size / Records 是否明显大于平均值如果是 shuffle read 倾斜可以看慢 task 所属分区对应的 key 分布情况对 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.instances50只是“最多申请 50 个”并不代表一定立刻拿到 50 个。如果队列资源不足Spark 只能等待。此时先看yarn application -list或对应调度平台确认应用实际获得了多少容器。必要时减少每个 executor 的 cores/memory让单位资源更容易被调度或者把大 Executor 拆成多个小 Executor。如果集群长期混部多个 Spark 应用建议设置spark.scheduler.modeFAIR让不同 job 的资源竞争更公平避免某个大应用长期占满队列影响其他作业的并行度。最后分享一个我自己的习惯每次写新的 Spark 作业启动时先把默认并行度和首份数据的分区数打出来。print(defaultParallelism:, sc.defaultParallelism) print(numPartitions:, rdd.getNumPartitions())跑完第一个 stage 后看一眼 Spark UI再决定要不要调repartition或coalesce。真正把“分区数”和“并行度”分开理解后排错和调优会顺畅很多。以上内容来自我在真实集群上踩过的坑不一定对每个业务都最优但覆盖了绝大多数场景的判断思路。
返回列表