ARTICLE DETAIL

资讯详情

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

Spark分区与并行度全解析:从RDD到Shuffle的调优指南

Spark分区与并行度全解析:从RDD到Shuffle的调优指南 写这篇的起因是被一个朋友问到一个看似特别基础的问题Spark 的 partition 和并行度到底由什么决定。这问题看着简单真要说明白能把 RDD、Stage、Shuffle、资源调度全部串起来。我见过不少人栽在根上有人觉得分区数越大越好结果整个 Stage 变成几千个小任务在集群上空转有人读了一堆小文件每个文件一个分区调度开销比计算本身还大还有人搞不清楚 Spark SQL 的 200 和spark.default.parallelism到底谁说了算。这篇文章就把这条线彻底拆开适合准备面试、正在调优、或者刚学 Spark 想少踩坑的同学。1. 先把概念捋清楚分区、Task 与并行度不是一个东西1.1 分区是数据切片不是任务RDD 的全称是弹性分布式数据集抽象上可以把它理解成一份数据被切成若干份每一份就是一个 partition。数据分布式地存在集群多个节点上而计算的时候最基本的数据单位就是这个 partition——Spark 不会对整个 RDD 做计算而是逐个 partition 地处理。每个分区在物理上通常对应一段数据可能来自 HDFS 的一个 block、Kafka 的一个分片、或者内存集合里的一份切片。举个例子树上结了 100 个苹果你要统计这批苹果的品质不可能让一个人把 100 个苹果全抱在怀里慢慢数那样又慢又累。更合理的做法是把苹果按篮子分成 10 份10 个人每人拿一篮去处理。这 10 个篮子就是 10 个 partition苹果总量还是 100但处理粒度变成了 10 份。RDD 就是这样的逻辑分区的数量是你可以观察、可以调整的rdd.getNumPartitions()一行代码就能看到当前有多少个分区。1.2 Task 才是真正跑起来的东西分区是数据的视角Task 是计算的视角。Spark 把一个作业提交后最终落在 Executor 上执行的最小单元叫 Task每个 Task 负责某个 Stage 里、某个分区上的数据处理逻辑。这里有一个非常关键的关系一个分区对应生成一个 Task正常情况下两者是一比一的关系。也就是说某个 Stage 里的 Task 总数基本就等于该 Stage 最后一个 RDD 的分区数。继续用苹果举例子10 个篮子就要安排 10 个 Task 去处理1 个 Task 对应 1 个篮子。但 Task 并不是一次全部在跑要看你给了多少工作线程这个稍后再细说。不过有一条要建立起来如果你发现一个 Stage 里生成了 3000 个 Task那背后一定有一个 RDD 被切成了 3000 个分区。往前查总能找到是哪一步输入、哪个算子造成了这个分区规模。这是排查问题最常用的思路。1.3 “并行度”被滥用的原因“并行度”是这三个概念里最容易混淆的。日常大家说“并行度取决于分区数”严格讲不完全准确。Task 是并发调度的单位但 Spark 是线程模型一个 Executor 能并发的 Task 数受它拥有的 CPU 核数也就是spark.executor.cores限制。因此一个 Stage 真正“同时执行”的 Task 数上限大约等于整个集群分配给这个应用的 CPU 核数之和。那为什么大家都习惯把分区数等同于并行度因为分区数决定了 Task 总数也就决定了这个 Stage 理论上能达到的最大并发能力。假设你有 20 个 Executor、每个 4 核一共 80 个并发槽位分区数只有 10那最多同时跑 10 个 Task剩下 70 个槽位空转这时候提升并行度最直接的办法就是增加分区数。反过来如果你有 1000 个分区一批也只能跑 80 个剩下的 920 个排队等待。用食堂打饭来类比你一下就记住了分区是盛菜的盘子数Task 是每个窗口处理这份菜的动作并行度是食堂同时能开几个窗口。盘子多说明要做的事多但同一时刻能出几个菜取决于窗口数——也就是 CPU 核数。调优的本质就是让“盘子数”和“窗口数”匹配别让任何一边闲着。2. 分区数到底由谁决定先看数据从哪来2.1 读文件时RDD API 和 Spark SQL 的行为不一样最常见的输入是从 HDFS 读文件。这里要区分两件事RDD API 和 Spark SQL 的分区逻辑并不完全一样别混在一起记。如果用 SQL 或 DataFrame 读Spark 会参考spark.sql.files.maxPartitionBytes默认 128MB同时参考spark.sql.files.openCostInBytes默认 4MB。它会尽量把物理文件按大小切成接近 128MB 的逻辑块小文件会被合并。所以同样一堆小文件用 DataFrame 读出来的分区数可能远小于文件个数反过来一个 10GB 的大文件大概率被切成 80 个左右的分区。这里的 128MB 和 HDFS block size 不是同一个东西这是 Spark 读文件时的“虚拟分片”上限。如果走老式 RDD API比如sc.textFile那走的还是 Hadoop InputFormat 的 split 逻辑基本一个 HDFS block 对应一个分区。一个 1GB、block 128MB 的文件分区数一般就是 8 个左右。这里有三个高频坑第一文件数多但每个都很小几百 KB 的日志RDD API 下每个小文件至少一个分区1000 个小文件就是 1000 个分区每个 Task 只处理一点点数据调度开销占了大头。第二如果文件是 gzip 这种不可切分的压缩格式哪怕文件有 2GB它也只能作为一个分区整体读取读入后的并行能力被压死在这个文件上。第三minPartitions参数不是你想的那样。很多教程说textFile默认分区数是 2那其实只是个兜底值真正读 HDFS 文件时分区数由 split 数决定并不会真的只有 2 个。这个参数只在你手动指定一个更大下限时才有意义。2.2 集合输入时parallelize 的默认值很“随缘”sc.parallelize(collection)是本地开发和测试最常用的入口它的分区逻辑也很容易忽略。不指定第二个参数numSlices时分区数取SparkContext.defaultParallelism而这个默认值在 local 模式下是当前机器的 CPU 核数——如果你用local[*]那就是宿主机所有逻辑核在 Standalone 或 YARN 模式下如果没配spark.default.parallelism一般取集群分配给应用的总核数保底为 2。开发机上常见的问题就是这样来的一台笔记本 16 核parallelize出来的 RDD 默认就是 16 个分区测试数据量小每个分区几千条完全没问题。但代码一上集群没显式配置的情况下集群总核数可能是 500分区数直接变成 500每个分区数据量反而变小Task 数量暴涨。同一个作业两种环境分区数完全不同行为差异就是这么来的。所以只要是用parallelize这种入口我都建议第二个参数显式传别赌默认值。2.3 Kafka 输入时分区数被源头锁死流式场景下Kafka partition 是并行度的第一约束条件。Spark 启动后默认每个 Kafka partition 分配一个读取任务去消费所以你的 topic 有 20 个 partition输入流就有 20 个分区。这里没法靠 Spark 端硬加并行度因为源头就 20 个分片下游就算 repartition 成 100 个也没法增加源头读取的并行度。想提升整个链路吞吐先看 Kafka 端的 partition 数够不够再看下游处理是否需要 repartition 来放大并行。这个约束很多人一开始没意识到等发现消费速度上不去时才回头看发现卡在源头。3. 转换算子怎么改分区数很多人栽在这里3.1 窄依赖算子分区数不变别指望它提并行RDD 经过转换算子后分区数不是一成不变的但也不是所有算子都会变。按依赖类型来分就很清楚窄依赖算子比如map、filter、flatMap、mapPartitions每个输入分区只产出一个输出分区分区数保持不变只是分区内的数据被筛选或变换。你可以把它理解成流水线上给每个产品贴标签流水线本身是几条就是几条不会变多。所以如果你的 RDD 一开始只有 4 个分区后面连刷十个map分区数还是 4。有人天真地以为多写几个算子并行度就上去了这是误解。想提高并行度要么在源头把数据切细要么显式做重分区。3.2 宽依赖算子shuffle 分区数有默认规则但不一定合理真正让分区数发生大变化的是宽依赖也就是会触发 shuffle 的算子reduceByKey、groupByKey、join、distinct、repartition等。shuffle 意味着数据要按 key 重新路由到下游分区这个下游分区数就是新 RDD 的分区数它不由上游直接决定而是由分区器和参数决定。重点来了很多人写reduceByKey(func)不传分区数以为它会沿用上游其实内部会调用defaultPartitioner优先级大致是父 RDD 若有现成的 Partitioner 且分区数较大就沿用否则如果设置了spark.default.parallelism就用它再兜底到父 RDD 里的最大分区数。这个默认逻辑会导致一个非常常见的现象从 HDFS 读出来的 RDD 本来有 800 个分区做了一次reduceByKey之后还是 800 个分区每个分区数据量小、Task 又多。如果你其实只需要 100 个分区做聚合那凭空多出来的 700 个 Task 都在浪费资源。反过来上游只有 4 个分区的小测试数据reducer 也只生成 4 个分区上集群之后根本撑不起并行。所以宽依赖算子我强烈建议显式传numPartitionsreduceByKey(func, numPartitions)、groupByKey(numPartitions)让分区数掌握在自己手里。提示简单判断一个算子会不会改分区数就看它是否触发 shuffle。窄依赖算子不改分区数宽依赖算子会按分区器重新切分。这条规律比死记算子列表靠谱得多。3.3 主动重分区repartition 和 coalesce 该怎么选除了默认行为Spark 给了两个主动调整的工具repartition和coalesce。这两者的区别很经典也是面试高频题。repartition(n)本质是coalesce(n, shuffletrue)数据会按 hash 打散到 n 个分区可以增加也可以减少分区代价是全过程有 shuffle网络和磁盘开销不小。coalesce(n)默认不 shuffle只是把上游多个分区的数据直接合并到下游少数分区适合大幅减少分区数比如最后写文件时想控制文件个数。但要注意从 2000 个分区直接coalesce到 2 个会发生数据堆挤因为下游只有两个 Task所有数据都往两个桶里汇聚某个 Executor 可能瞬间内存爆炸。这时候要么用coalesce(n, shuffletrue)要么直接用repartition(n)重新打散一次。还有一个partitionBy算子常见于写数仓时决定输出文件怎么分桶。它给 RDD 打上 Partitioner让后续依赖它的 Stage 在 shuffle 时知道怎么路由写中间结果时很实用能减少下游重复计算。对比点repartition(n)coalesce(n)是否 shuffle有默认无shufflefalse分区增减可增可减通常只减少数据均衡性按 hash 打散总体较均匀大幅减少时可能严重不均衡适用场景提高并行度、重分布数据输出前减少分区、控制文件数4. 并行度的真正上限Stage 划分与资源槽4.1 Task 数从来不是针对整个作业说的作业提交后Spark 会构建 DAG然后以宽依赖shuffle为界把 DAG 切成多个 Stage。为什么要讲这个因为在 Spark 里“Task 数”从来不是对一个作业说的而是对每个 Stage 说的。作业里第一个 Stage 的 Task 数量由输入 RDD 的分区数决定中间某个 Stage 由它自己的最后一个 RDD 决定shuffle 之后的 Stage 则由上一阶段写入的分区数决定。所以你去 Spark UI 看同一个 job 里不同 Stage 的 Task 数可能差异巨大前面读文件 100 个 Task后面 join 之后变成 800 个 Task这都正常。排查性能问题时要按 Stage 逐个看不能拿整个作业的平均数来说事。如果一个作业慢先定位是哪个 Stage 慢再去看那个 Stage 的 Task 数合不合理。4.2 并发执行上限Executor 核数才是真正的天花板Task 排队执行靠的是 Executor 里的线程池。一个 Executor 申请了spark.executor.cores4默认最多同时执行 4 个 Task如果设了spark.task.cpus大于 1会按倍数扣。假设集群给你 20 个 Executor那这个应用同时能跑的 Task 数上限就是 80。这个数才是真正的并发执行上限也就是资源层面的并行度天花板。很多面试题问“Spark 的并行度取决于什么”比较完整的回答框架应该是从数据层面看取决于各 Stage 中 RDD 的分区数上游由数据源和分片逻辑决定下游由 shuffle 分区参数决定从资源层面看取决于 Executor 数量乘以单个 Executor 的核数。整体并行度是这两者的交集。这样答比单纯说“取决于分区数”完整得多也更能体现对调度的理解。4.3 分区数、并行度与数据量的三角平衡那分区是不是越多越好不是。每个 Task 都有调度、序列化、任务描述的开销成百上千个只处理几 KB 数据的 Task纯粹是在交调度税。反过来分区太少也不行单 Task 处理数据量过大会导致 GC 频繁、内存溢出、磁盘 spill 变多。实践中看两个参考。一个是总分区数保持在集群总核数的 2 到 3 倍给执行完毕的快慢差异留缓冲因为有些 Task 跑得快有些跑得慢多一点分区能让先跑完的 Executor 接上后续任务避免尾巴拖慢整个 Stage。另一个是单个 Task 的输入数据量控制在 100MB 到 1GB 之间具体看计算复杂度聚合类、排序类要偏保守如果单 Task 数据量上 2GB 了就要赶紧调大分区数。这两个经验值都不需要特别精确但作为起步判断非常管用。5. 参数清单与经验法则照着抄就能少踩坑5.1 两个核心参数defaultParallelism 和 sql.shuffle.partitions调分区绕不开两个参数spark.default.parallelism和spark.sql.shuffle.partitions。前者是 RDD API 的默认分区来源影响parallelize的默认切片数、宽依赖不显式设置分区时的兜底值后者是 Spark SQL 里所有 join、group by、distinct 等 shuffle 操作的下游分区数默认 200。两者互不替代。大多数人踩坑就踩在这里在 Spark SQL 里改了spark.default.parallelism结果发现 group by 还是 200 个分区因为 SQL 默认只认spark.sql.shuffle.partitions。你要记住一条边界凡是走 SQL 或 DataFrame 的 shuffle看spark.sql.shuffle.partitions凡是走 RDD 算子且不显式传分区数看spark.default.parallelism。两个都设数值尽量保持大致一致能少很多争议。5.2 一个可以直接抄的提交模板以一个 20 Executor × 4 核的 YARN 应用为例我常用的起点配置是这样总核数 80目标分区数 2403 倍。在spark-submit里我一般这么写spark-submit \ --master yarn \ --num-executors 20 \ --executor-cores 4 \ --executor-memory 8g \ --conf spark.default.parallelism240 \ --conf spark.sql.shuffle.partitions240 \ --conf spark.sql.files.maxPartitionBytes134217728 \ --conf spark.sql.files.openCostInBytes4194304 \ --conf spark.sql.adaptive.enabledtrue \ app.jar代码里也可以动态拿到核心数来算spark.sparkContext.defaultParallelism * 3但要注意此时defaultParallelism可能已经被配置值改掉了别把 240 又乘了 3。所以更推荐的做法是外部传参、代码里读取。如果用的是 DataFrame 做聚合建议代码里也显式指定一下 shuffle 分区比如设置 SQL 参数或df.repartition(240)这样即使提交脚本忘了传代码本身也有个保底值。数据量级小时200 的默认值其实够用没必要一上来就堆几百上千个分区。5.3 经验法则读阶段、写阶段、观察指标我再给几条自己实测下来比较稳的法则。读文件阶段理想分区数约等于文件总大小除以 128MB但不要超过输入文件数的 1.5 倍小文件场景除外。shuffle 阶段分区数取总核数 2 到 3 倍宁多勿少。多了后面的 AQE 可能帮你合并少了在低版本 Spark 里是真刀真枪的性能损失。写文件场景最终输出分区数决定文件数写 HDFS 时不要故意把分区拉得很高控制文件个数比无限并行更重要。动态观察时Spark UI 里看 Shuffle Read Size 和 Records。如果单 Task 读入数据超过 1 到 2GB说明分区太少看 Scheduler Delay 和单 Task 耗时如果 Task 数量几千个但每个都只跑几十毫秒说明分区太多。5.4 Spark 3 的 AQE让分区数“自动收敛”如果你们用的是 Spark 3 以上强烈建议开spark.sql.adaptive.enabledtrue3.0 需要手动开启3.2 之后的大多数发行版已经默认开启。AQE 里有三个能力和分区直接相关。第一动态合并 shuffle 分区它会在运行时根据实际 shuffle 数据量把过小的分区合并掉避免小 Task 满天飞。第二倾斜 join 优化它会识别出数据倾斜的 join 键把大分区拆细。第三动态切换 join 策略它可以根据统计信息把排序合并连接改成广播连接。开了 AQE 之后你甚至可以把初始分区数略微调大让运行时自己收拢容错率高出不少。但注意AQE 只对 shuffle 之后生效读文件阶段那几百个小文件它管不了小文件问题还是要在源头解决。6. 排查实录几个踩过的坑和解决方案6.1 小文件把分区数撑爆了有次排查一个离线任务现象是每个 Stage 都有几千个 Task但每个 Task 只处理几 KB 数据整体跑得又慢又飘。查源头上游把一天日志按小时落了几百个文件每个文件才几百 KB用老 RDD API 读进来每个文件一个分区。解决方案分两步读的时候改用 Spark SQL 读配合openCostInBytes4MB让这些小文件被合并成大逻辑分片如果必须用 RDD API可以用newAPIHadoopFile指定CombineTextInputFormat把多个小文件合并成一个输入分片。改完之后分区数从几百掉到二三十整个作业时间缩短了一半多。这个问题在日志类数据场景特别普遍源头治理比下游反复 rebalance 有效得多。6.2 shuffle 后数据倾斜靠 repartition 是救不了的一个 join 作业两个流各自 200 个分区join 之后明显有个别 Task 的 Shuffle Read Size 高达几十 GB其他 Task 只有几百 MB。默认 hash 分区器按 key 的 hashCode 取模热点 key 一多数据必然集中在某个桶。最简单粗暴的办法是给热点 key 加随机后缀把一个大 key 拆成几十个小 key两阶段聚合后再合并如果 join 两边都是大表可以在加盐后改成两次 join 方案。业务允许的情况下也可以直接开 AQE 的spark.sql.adaptive.skewJoin.enabled它会自动把倾斜分区按比例拆小。这里的关键是别指望repartition解决倾斜——repartition只保证分区数不保证每个分区里数据量均衡。热点 key 还是那个热点 key换个桶照样挤。6.3 200 个分区把内存跑爆了另一个典型的 SQL 场景某张大表 group by 之后要排序输出数据量 10TB默认spark.sql.shuffle.partitions200每个分区要处理 50GB单 Task 内存根本不够shuffle 数据不断 spill 到磁盘然后又疯狂读盘速度惨不忍睹。我把参数调到 2000单 Task 数据量降到 5GB配合更大内存才缓过来。这类问题在 Spark UI 上很好认——查看某个 Stage 的 Shuffle Read Size如果平均值超过 1GB 且 Executor 出现 GC 或 OOM 报错第一步就是加大 shuffle 分区数别先忙着加 Executor 内存。加大内存也有效但治标不治本单点负担还是太重分区数才是把压力分散开的核心手段。6.4 本地没毛病上集群就“水土不服”这个坑几乎每个把本地开发代码丢到集群上的人都会遇到。本地local[8]时parallelize默认 8 个分区跑得欢快。同一个 JAR 上了 80 核的集群没设spark.default.parallelism读同一份小测试文件时分区数就飘到 80 甚至 160Task 全是调度开销速度反而更差。我的习惯是所有涉及并行度的关键入口都显式传参parallelize给numSlicesshuffle 类算子给numPartitions提交脚本固定spark.sql.shuffle.partitions和spark.default.parallelism。让代码在本地、测试、生产三个环境行为一致排查问题时就不会被“环境差异”这个变量干扰。测试环境的数据量再小分区策略也要和生产对齐否则你测的只是“本地行为”不是“生产行为”。最后聊一点个人体会分区和并行度这两件事本质上是在调“数据切分的粒度”和“计算资源的匹配度”。很多调优看似玄学其实只要顺着一条线去查——某个 Stage 的 Task 数是否和资源槽匹配、单 Task 数据量是否超出合理区间、shuffle 分区参数有没有被默认值绑架——基本都能定位到根因。Spark UI 上的每个 Stage 都是一份体检报告分区数、Shuffle Read、Task 耗时都在那写着比猜来猜去靠谱得多。你每调一次分区都建议把修改前后的 UI 指标截图留下来时间长了就是自己的一套调参经验库。
返回列表