☰
Spark合并参数辨析:coalesce、repartition与AQE实战指南
2026/9/30 20:54:59 网站建设 项目流程

上周有个同事跑过来问我:“我把任务最后一步的repartition(200)改成了coalesce(200),怎么反而变慢了?”我一看他改的那段代码,源表数据总共就只有 8 个分区,写表前想用 coalesce 把输出文件数量“合并”一下,结果数据全被压到了少数几个执行器上,文件没少,任务倒是慢了两倍。

这个案例几乎是 Spark 合并参数误用的标准模板。平时大家嘴上说的“合并”,在 Spark 里其实是一族完全不同的东西:有只改物理分区的coalesce,有走全量 shuffle 的repartition,有写文件时控制目录的partitionBy,还有从 Spark 3.0 开始默认开启的自适应查询执行(AQE)自动分区合并。所有这些都被笼统地叫“合并参数”,但底层逻辑、触发条件、适用场景区别极大,用错了不是慢一点,而是直接废掉整个任务的并行度。

这篇文章我想把这些容易搞混的算子、参数和机制理一遍,结合我实际跑任务时踩过的坑,说清楚它们各自在合并什么、什么时候能混用、什么时候绝对不能混用。适合正在做 Spark 数据分析、数仓清洗链路,以及被小文件问题折磨过的朋友参考。

1. 先从最容易被误用的两个算子说起:coalesce 和 repartition 的角色边界

1.1 为什么一句“把分区合并一下”其实是两套逻辑

coalesce(n)和repartition(n)从 DataFrame API 看都返回一个新的分区数为 n 的 DataFrame,这也是很多人觉得“这俩不是一样吗”的原因。但去看物理执行计划,完全是两条路。

repartition(n)会触发一次真正的 shuffle,也就是把所有数据按分区编号重新打散,每个分区里的数据经过网络传输到对应的新分区。这个过程是宽依赖,开销明显,但是它能把数据均匀地重新分布。coalesce(n)默认不触发 shuffle,它做的事情只是把多个物理分区“搭桥”到一起,例如把原来的 10 个分区合并成 3 个分区,本质是建立窄依赖,不需要网络传输。

打个比方:repartition 相当于把仓库里所有货架上的货全部倒出来,按照新的货架编号重新分配;coalesce 则是把相邻的几个货架直接拼在一起,货还是原来的货,只是架子变少了。听起来 coalesce 更省事对不对?对,但“拼架子”这个动作有一个天然限制——coalesce 只能减少分区数,不能增加。你数据本来 8 个分区,coalesce(200)执行完最多还是 8 个分区,多出来的分区编号是空的。

这是第一层区分的关键:想增加分区数来提升并行度,只能 repartition;想减少分区数且数据量不太离谱,可以用 coalesce 省掉一次 shuffle。但省掉 shuffle 不等于没有代价,下面展开说。

1.2 什么时候该用 coalesce,什么时候必须用 repartition

我自己的选型习惯可以整理成一张表,日常基本够用:

场景推荐方案原因
分区数从几千降到几百,数据分布均匀coalesce(n)无 shuffle,速度快
分区数很少,但想增加并行度repartition(n)只有 shuffle 才能把数据打散
上游某个 key 数据量极大,存在倾斜repartition(n, col)或加盐方案只减少分区不能解决倾斜,可能更糟
写少量输出文件,数据总量很小先repartition(1)或coalesce(1),但注意内存单分区写文件简单,但数据大时要谨慎
合并 ratio 很小(比如 100->1)尽量repartition(1)或调整策略coalesce 可能把一个执行器压爆
最后写表要精确控制文件数量repartition(n)更稳输出文件由 shuffle 后分区数量决定,repartition 语义最明确

关于第 4 行和第 6 行要特别说明:很多人看到 coalesce 省 shuffle,于是所有需要减少分区的地方都无脑 coalesce。但当分区数缩减比例很大时,coalesce 会把原本分布在几十个 task 里的数据合并给极少几个 task 处理。比如你数据总量 1TB,从 1000 个分区coalesce(10),每个 task 要处理 100GB 数据,这时候网络传输倒是省了,但写文件、序列化、内存压力全集中在这 10 个任务上,最后可能 OOM 或者长时间卡在最后几个 task 上。

所以我的经验是:分区规模比较小(比如几百到几十)或者数据总量不大,coalesce 很香;如果是从很大的分区数压到很小的分区数,宁可走一次 repartition,让每个执行器都分摊一点数据,至少任务能稳定跑完。

1.3 coalesce 的隐藏成本:没有 shuffle 不代表没有开销

再深入一个容易忽略的细节:coalesce 建立的是窄依赖,但窄依赖意味着同一个阶段里一个 task 可能要处理多个原始分区的数据。在 Spark 的调度模型里,一个 stage 的 task 数量取决于该 stage 最后一个 RDD 的分区数,所以 coalesce 完之后 task 数量确实变少了。问题是,task 变少并不等于每个 task 的工作量合理。

举个实际例子:两个原始分区一个 10GB,一个 10MB,coalesce 到 1 个分区后,单个 task 要处理 10.01GB 的数据。如果这时候集群有 20 个执行器,其他 19 个执行器只能闲着,因为它们对应的分区都已经被合并掉了。这种情况 Spark UI 的 Stage 页面看起来特别典型:大量 task 秒完,最后剩下一个 task 跑几个钟头。

还有一点是关于空分区的。coalesce(n)在分区数减少时可能产生空分区,这些空分区虽然不会产生实际 task,但会影响下游算子对分区数的预期。比如你 coalesce 之后做mapPartitions并且内部有初始化资源、连接池之类的逻辑,空分区依然要走一遍处理逻辑,白白耗时。这种问题排查起来很隐蔽,因为执行计划看不出任何异常。

2. 写文件场景中,partitionBy、bucketBy、coalesce 到底在合并什么

2.1 partitionBy 合并的是目录,不是数据分区

很多人在写 Hive 表或者 Parquet 文件时,看到输出目录下有几万个文件,第一反应就是“我给分区字段加一个 partitionBy 是不是就合并了”。这是把两个完全不同层面的概念混在了一起。

partitionBy("dt")控制的是目标表的目录结构,写入后数据会按照 dt 的不同取值拆到不同目录下,比如dt=2024-01-01、dt=2024-01-02各一个目录。它做的事是“逻辑分区字段映射到物理目录”,不是把某个目录内部的碎片文件合并。

我见过最典型的误用是:一张表按城市字段动态分区,任务跑完以后每个城市分区目录下都有几十个几十 KB 的小文件。开发以为是 partitionBy 参数没设好,其实 partitionBy 根本没有“合并文件”的职责。真正决定文件数量的,是往每个目标分区写数据时经过了几个输出 task,以及是否触发了动态分区写入的特殊路径。

所以正确理解是:partitionBy负责把数据分门别类放进不同目录;coalesce/repartition负责控制每个目录里由多少 task 写出的文件数。这两个参数配合用,才能既保证目录结构正确,又保证每个目录下的文件数量可控。

2.2 bucketBy 和 repartition 的绑定关系容易被忽略

如果说 partitionBy 还有很多人认得出是“目录分区”,那bucketBy就是更冷门也更容易埋雷的一个。bucketBy(n, "id")写在 saveAsTable 前面时,它声明的是这个表在元数据里被分成 n 个 bucket,底层文件的组织方式应该和 bucket 一一对应。但关键问题在于:bucketBy 这个声明本身并不会自动触发对应的分区逻辑。

正确的写法通常长这样:

df .repartition(10, col("id")) .write .bucketBy(10, "id") .sortBy("id") .format("parquet") .saveAsTable("user_bucket")

这里的repartition(10, col("id"))保证了写表前的数据经过一次按 id 哈希的 shuffle,输出分区数刚好等于 bucket 数。如果不加这一步,Spark 每个写 task 都会生成一个文件,最后落盘文件数量和 task 数量挂钩,跟 bucket 数量完全不匹配,查 meta 信息时表是 10 个 bucket 的元数据,实际目录里却躺着几百个文件,bucket pruning 性能优势直接归零。

bucketBy还有一个常见误区:它只是元数据层的组织策略,不等于“按 bucket 去重”或者“减少数据量”。它的作用是让相同 id 的数据落在同一个 bucket 文件里,方便后续 Join 时做 bucket 裁剪,减少扫描量。它和“合并文件”“减少文件数量”之间的关系是间接的:只要 repartition 的个数匹配 bucket 个数,那么一个 bucket 就对应一个文件,这确实能达到物理上减少文件数的效果,但前提是你真的做了 repartition。

2.3 maxRecordsPerFile 才是控制单文件大小的另一个闸门

分区数和文件数并不是完全的一一对应关系,中间还夹着一个spark.sql.files.maxRecordsPerFile。这个参数的默认值是 0,表示不做限制;一旦设置成某个正整数,比如 100 万,那么单个文件里记录数超过 100 万时,Spark 会继续写第二个文件。

这解释了为什么有些时候你明明已经把分区数压到 1 了,输出目录里还是有两个文件——因为一次写入的记录数超过了阈值,触发了文件拆分。反过来,如果你希望尽量让每个文件都大一点、数量少一点,可以把这个参数调大或者保持 0。

不过这里有个合力问题:maxRecordsPerFile控制的是“每个文件最多多少条记录”,coalesce/repartition控制的是“写数据时有多少个输出分区”,两者叠加才决定最终文件数。只调 repartition 不调 maxRecordsPerFile,遇到超大数据量的分区还是会被拆;只调 maxRecordsPerFile 不调分区数,文件总量可能依然很大。实际调优时先确定目标文件大小,再倒推分区数和记录数,比单点调一个参数靠谱得多。

3. AQE 时代的自动合并参数:手动合并和自动合并的分工

3.1 spark.sql.shuffle.partitions 在开启 AQE 前后的含义变化

spark.sql.shuffle.partitions是 SQL 作业里 shuffle 阶段 reducer 的默认分区数,默认值 200。以前调优这个参数几乎是分析型任务的必修课:数据量大就调大,数据量小就调小,调不好不是小文件爆炸就是并行度不足。

但 Spark 3.0 之后 AQE 默认开启,情况发生了变化。开启 AQE 后,spark.sql.shuffle.partitions的角色更接近“初始分区数”,而不是最终分区数。AQE 在 shuffle 结束后会统计每个分区的数据量,根据设定的目标分区大小自动把相邻的小分区合并起来,最终输出分区数可能远小于 200,也可能在数据量极大的时候大于 200,取决于spark.sql.adaptive.advisoryPartitionSizeInBytes这个目标值。

所以现在很多任务不再需要手动去反复试spark.sql.shuffle.partitions了,只要你把 AQE 的一些参数设好,Spark 自己会在每个 shuffle 边界上动态调整分区数。当然,前提是你没有在代码里显式调用repartition(n),因为显式调用的分区数是硬指定,不会因为 AQE 开启就自动改。

3.2 一组可以照着抄的 AQE 合并参数配置

下面这组配置是我在 100 到 500 个执行器的集群上日常使用的,比较中庸,适合大多数离线分析场景:

spark.sql.adaptive.enabled=true spark.sql.adaptive.coalescePartitions.enabled=true spark.sql.adaptive.coalescePartitions.minPartitionNum=16 spark.sql.adaptive.coalescePartitions.initialPartitionNum=200 spark.sql.adaptive.advisoryPartitionSizeInBytes=134217728

逐条解释一下:

  • spark.sql.adaptive.coalescePartitions.enabled:是否允许 AQE 自动合并 shuffle 分区,默认就是 true,一般不用动。
  • spark.sql.adaptive.coalescePartitions.minPartitionNum:合并后的分区数下限,默认 1。我一般调成 16 到 32,避免极端情况下自动合并把所有分区压成一个,导致并行度归零。
  • spark.sql.adaptive.coalescePartitions.initialPartitionNum:第一次 shuffle 时的初始分区数,如果不设置就取spark.sql.shuffle.partitions。数据量差异特别大的库,可以单独设成几百。
  • spark.sql.adaptive.advisoryPartitionSizeInBytes:AQE 希望每个 shuffle 分区达到的目标大小,默认 64MB。如果下游任务每个分区处理逻辑很重,可以调大到 128MB 或 256MB,减少分区数、降低调度开销;如果单个分区处理逻辑很轻,可以往 32MB 调,提高并行度。

这套配置最舒服的地方在于:不同表、不同查询的数据量差异很大时,你不再需要为每个任务手动猜分区数。比如一张表清洗后只有 2GB,初始 200 个分区,每个分区只有 10MB,AQE 会自动合并成 16 到 32 个分区左右;而另一张 500GB 的表,初始 200 个分区,每个分区约 2.5GB,AQE 不会强行合并,甚至会通过更多初始分区来分摊。这比手工调参区间宽多了。

3.3 自适应合并什么时候不会触发:动态分区写入的坑

AQE 自动合并虽然有诸多好处,但它不是万能的。我遇到过最多的问题是:AQE 开得好好的,写表时还是生成了大量小文件,最后查执行计划发现动态分区写入路径下,合并并没有按照预期发生在每个目标分区上。

动态分区写入指的是INSERT OVERWRITE TABLE t PARTITION(dt) SELECT ...这种写法,目标分区字段的值由数据本身决定。Spark 在写入时会为每个目标分区单独分配 writer,如果上游数据经过 shuffle 后某个分区的数据被拆到多个 reducer 文件片段,最后写同一个目标分区时会有多个文件。AQE 的合并作用在 shuffle 输出的 reducer 分区层面,但动态分区写入的“每个目标分区多个文件片段”问题,往往由写入机制本身造成,不是说 AQE 完全管不了,而是它合并的是“按哈希分区后的 reduce 分区”,不是“Hive 表的物理分区目录”。

这种情况下更有效的做法是手动控制数据分布。比如用 SQL 的DISTRIBUTE BY dt:

INSERT OVERWRITE TABLE t PARTITION(dt) SELECT * FROM src DISTRIBUTE BY dt

DISTRIBUTE BY dt会让同一个 dt 值的数据在 shuffle 后落入同一个 reducer 分区,这样每个目标分区只对应一个输出文件,文件数量直接由 dt 的枚举值数量决定。如果某些 dt 数据量特别大,一个文件如果撑不住,可以改成DISTRIBUTE BY dt, rand()之类的方式,让同一个 dt 内部进一步拆分成均匀粒度,但这样又会多文件。所以这个 trade-off 需要结合数据量级来判断。

4. 小文件治理实战:把参数组合起来而不是一个一个试

4.1 先定位小文件是在读还是在写阶段产生

小文件问题几乎每个 Spark 任务都会碰到,但很多人一上来就改参数,忽略了定位阶段。同样的“文件巨多”,读阶段和写阶段的解法完全相反:读阶段希望把多个小文件打包成一个分区 Task 去处理,减少 task 数量;写阶段希望输出 task 少一点、每个文件大一点,减少落盘文件数。

先用命令把现状统计清楚:

hdfs dfs -count -v /warehouse/ods_order hdfs dfs -ls /warehouse/ods_order/dt=2024-01-01/ | wc -l

如果输入目录有 1 万个小于 1MB 的文件,问题大概率在读阶段,需要调spark.sql.files.maxPartitionBytes和spark.sql.files.openCostInBytes;如果输入文件都很正常,但写完后输出目录文件巨多,问题在写阶段,需要调分区数和动态分区写入策略。跳过定位直接调参,经常是白忙活。

4.2 读阶段合小文件的参数:maxPartitionBytes 与 openCostInBytes

这两个参数经常被混淆,我放在一起说。

spark.sql.files.maxPartitionBytes控制的是读文件时单个分区最多能读多少字节,默认 128MB。比如设置成 256MB,那么一个分区最多承载 256MB 的数据扫描量,超过的部分会分给下一个分区。它限制的是“分区最大容量”,不是“最小容量”。

spark.sql.files.openCostInBytes则是给每个小文件加的“打开成本”权重,默认 4MB。一个真实大小只有 100KB 的小文件,在 Spark 做分区规划时会被当成一个 4.1MB 左右的文件参与计算,多个小文件的成本累加起来超过分区上限后,就合并到同一个分区里处理。所以把spark.sql.files.openCostInBytes调大,例如 64MB,会显著减少读小文件时的分区数量,让更多小文件被同一个 Task 一起扫完。

典型的读阶段配置:

spark.sql.files.maxPartitionBytes=268435456 spark.sql.files.openCostInBytes=8388608

这里要提醒一句:这两个参数调大后 task 数量减少,单个 task 扫描的数据量变大,算子内部如果有关键耗时的初始化逻辑,依然会成为瓶颈,所以不能只靠参数硬扛,最好配合后续的重分区一起看效果。

4.3 写阶段推荐的组合套路与验证方式

写阶段的组合没有万能公式,但有三套基本功,按场景套用基本能解决 90% 的问题:

第一,非分区表和分区字段数据量均匀时,最简单粗暴的就是repartition(n),n 等于你想输出的文件数。这比coalesce(n)语义更明确,避免大比例合并且 task 饿死的问题。如果数据总量不大,例如 1GB,repartition(1)或coalesce(1)都行,我一般选repartition(1),图一个确定性。

第二,有动态分区的表,用DISTRIBUTE BY 分区字段或它的 DataFrame 等价写法repartition(分区字段)。这一种能保证同一个分区的数据尽量落到同一个 reducer,再配合sortBy(分区字段),写出来的文件在一个分区内还是有序的,查询性能更稳。

第三,数据量特别大、单分区文件还想要均匀大小的,用repartition(分区字段, 随机列)或者 SQL 里的DISTRIBUTE BY 分区字段, rand()。这样同一个分区会被拆成多个均匀的块,文件大小均衡,不会出现一个分区 10GB、另一个分区 10MB 的极端情况。

验证方式:跑完任务后看输出目录文件大小分布,最好写个小脚本做统计,别只看文件数量。文件数量正常但大小差异 1000 倍,后续读取依然难受,这往往是数据倾斜在写阶段的另一种表现。

4.4 一个典型的网约车数据清洗链路优化前后对比

之前做过一个网约车订单数据的离线清洗任务,场景很典型:每天从上游同步订单明细到 Hive,然后 Spark 做清洗,按城市分区写回结果表。优化前任务跑完,整张表文件数量有接近 3 万个,平均文件大小不到 10MB,每天跑批耗时 40 多分钟,下游分析任务读这张表也要 20 多分钟。

当时的执行计划里,清洗逻辑做完之后没有显式重分区,直接动态分区写表。Stage 里 reducer 数量是默认的 200,但城市分区有 30 多个,每个城市的数据又被若干 reducer 各自写出一份,所以小文件自然很多。

优化后的链路加了两个关键动作:写表前先repartition(col("city")),把城市分布规整到同一批 reducer;同时开启 AQE 并设置advisoryPartitionSizeInBytes=128MB,保证数据量大的城市不会因为单 reducer 处理吃力而整体卡住。最后再配合分区字段排序写表。

优化结果对比如下:

指标优化前优化后
结果表文件总数约 3 万个约 3000 个
平均文件大小小于 10MB约 80MB
单次清洗耗时40 分钟以上15 分钟左右
下游查询扫描耗时20 分钟以上5 到 8 分钟

这个案例的启发是:小文件治理不是靠某一个参数的神奇配置,而是把读阶段分区策略、写阶段重分区策略、AQE 自动合并三者放在同一条链路上统一设计。每层只解决自己那一层的文件粒度问题,最终文件数量才可能收敛到合理范围。

5. 排查合并效果时,我建议大家看的两个地方

5.1 执行计划里的 Exchange 节点决定合并是否真的发生

无论改什么合并参数,我建议第一步都是看物理执行计划。explain(true)输出的 plan 里最直观的标记就是Exchange节点。

df.coalesce(10).explain(true) df.repartition(10).explain(true)

repartition的物理计划里会出现Exchange hashpartitioning(...),说明这里有一次真正的 shuffle;coalesce默认情况下没有Exchange节点,只有Coalesce 10之类的逻辑节点。如果你打开 AQE,计划里还会出现AdaptiveSparkPlan,此时Exchange的数量、分区策略会由 AQE 在运行中调整,最终执行前你甚至可能看到分区数被动态修改。

这个习惯的价值在于:很多时候你以为自己设置了合并参数,但代码路径里压根没走到。比如某个转换算子内部先执行了 repartition,后来又有人加了 coalesce,两者相互抵消,最后执行计划里全是多余的 Exchange。先看计划,能少走很多弯路。

5.2 Spark UI 和线程监测工具辅助确认资源利用

执行计划看完之后,资源层面的确认同样重要。Spark UI 的 Stages 页面会直接显示每个 stage 的 task 数量、shuffle read/write 总量、GC 时间。如果某个 stage task 数量明显少于集群可用核心数,说明合并过头了;如果 task 数量很多但每个 task 处理的数据量只有几 MB,说明合并没生效,或者小文件读阶段没合并彻底。

比较隐蔽的问题是执行器线程层面的瓶颈。A 撤掉了高并行度后,某些 task 从表面看跑得很快,但 Executors 页观察内存和 GC 时发现老年代频繁回收,或者 CPU 利用不均衡。这时候我会抓一下执行器的线程栈,用jstack看主线程到底卡在 GC、等待网络还是 CPU 计算,顺便对照内存监测工具里各区域的使用趋势。很多“分区合并后变慢”的谜团,最后都能在 GC 时间或者序列化耗时上找到证据。

5.3 一个本质问题:合并分区不等于解决数据倾斜

最后说一个容易被合并参数掩盖的本质问题:数据倾斜。

我知道很多人遇到任务慢、某些 task 数据量超大时,第一反应是调大repartition分区数,或者把coalesce去掉。这些办法有时看起来有效,因为它把单个 task 的数据量分摊到了更多 task 上,但如果倾斜的根因是某个 key 占据了绝对大头,单纯调整分区数只是把一部分数据从一个大 task 挪到另一个大 task,倾斜依然存在。

更隐蔽的是 AQE 的自动合并也可能踩到倾斜。如果倾斜 key 的数据全部落入同一个 shuffle 分区,AQE 在合并时会把附近的小分区也拉进来,最后合并后的分区还是一个大胖子。处理倾斜要从数据本身下手,比如对热点 key 加盐做两阶段聚合、调整 Join 顺序、使用 broadcast join,或者开启spark.sql.adaptive.optimizeSkewedJoin.enabled让 AQE 在 Join 阶段拆分倾斜分区,而不是把所有希望寄托在“多合并几次”上面。

合并参数调整的是物理布局和数据分布的外壳,数据本身的分布特性才是内因。先把内因看清楚,再决定用哪些合并手段,才不会被任务表面跑通了但是依然慢的问题反复折磨。我个人的习惯是:每次改合并相关配置之前,先花十分钟看物理计划和文件分布统计,真的比随手调参数然后再试跑几个小时划算得多。

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

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

立即咨询