“我把一个 5000 分区的 DataFrame 执行了coalesce(10),写出去的还是 3000 多个小文件,这合理吗?” 上周在技术群里又看到类似的问题,说实话这类问题几乎每个月都会出现。Spark 里的合并参数看着简单,就coalesce和repartition两个方法,但底层的 shuffle 机制、数据分布逻辑、文件输出数量,每个环节都有不少坑。尤其在做了数据清洗、聚合分析这类任务后,分区数量失控导致小文件爆炸,或者数据倾斜导致某个任务跑几个小时,很多问题追根溯源都出在“合并分区”这一步。
我最早接触 Spark 时也在这上面栽过跟头。当时处理网约车订单数据,清洗完准备写回 Hive 表,随手用了repartition(1)想合并小文件,结果整个任务因为一次全量 shuffle 卡了快一个小时。后来才慢慢搞清楚:Spark 的“合并参数”从来不是简单的“把分区变少”,它背后是窄依赖与宽依赖的区别、shuffle 代价的取舍、文件输出数量的权衡,甚至还要考虑 AQE 动态合并的影响。这篇文章就把我这些年踩过的坑、总结的经验一次说清楚。
1. 合并参数到底是什么,为什么成为 Spark 调优的必修课
1.1 分区数量为何会失控
Spark 的并行度由分区(Partition)决定,每个分区对应一个 RDD 分区或 DataFrame 的一个 Task。默认情况下,读取 HDFS 上的文件时,一个文件块对应一个分区;但你做过 join、groupBy、distinct 这类宽依赖操作后,spark.sql.shuffle.partitions(默认 200)会决定结果分区数。问题就出在这:如果你读取了 5000 个小文件,又做了一次没有意义的 join,分区数就可能膨胀到几千甚至上万。
举例来说,我在做农产品价格数据的清洗时,源表是按天分区的,每天有 2000 多个小文件。任务先做了一次字符串替换、类型的转换,再做一个按品类维表的 join,跑完之后落地的文件依然有 2000 个。每个文件只有几十 KB,整个数仓被小文件塞满,查询时 NameNode 压力大,Spark 读取时 task 数过多,调度和序列化的开销远大于实际计算开销。
分区数量失控的直接后果,你可以类比成快递分拣:如果每个包裹都单独是一个袋子,分拣员(Executor)要把几万个袋子来回搬,搬运本身的时间就超过了实际分拣时间。合并分区的本质就是“把零散的包裹归拢到合适数量的袋子中”,让每个分区的数据量接近合理范围(通常一个分区 128MB~256MB 为佳)。
1.2 分区过多与过少的两难
分区太多,每个分区数据量太小,task 调度开销远大于计算开销;分区太少,数据量太大,单个 task 内存压力过重,可能出现 OOM、GC 频繁、甚至磁盘溢写。这里的“合适范围”没有绝对标准,但有一条经验法则:单个分区的数据量在 128MB~256MB 之间,消费者集群的并行度也能跟得上。
不少初学者以为合并参数就是把分区数调到“小”,觉得越小越好。其实不然。当你把分区合并到小于 Executor 总核数时,一部分核就空闲了,集群资源没有跑满。更关键的是,如果合并操作本身触发了 shuffle,数据落盘和网络传输的代价会远超你省下的那点 task 调度时间。合并参数的核心矛盾是:既要减少分区总数,又不能让合并动作本身变成性能瓶颈。
这就是为什么我要专门把coalesce和repartition拆开讲——它们一个不触发 shuffle,一个触发 shuffle,性能差异天壤之别,但很多人的选择是“随手点一个”。
2. 核心参数拆解:coalesce 与 repartition 的底层机制差异
2.1 窄依赖的天然优势:coalesce 为什么快
coalesce是 Spark 提供的分区合并方法,源码实现对应CoalesceExec。它的关键特征是不会触发强制 shuffle 重分区。当你调用coalesce(n)时,Spark 会尝试把多个现有分区合并到 n 个分区,尽可能保留数据在 Executor 上的位置,只在父 RDD 分区与子 RDD 分区之间建立窄依赖。
用生活化的比喻:coalesce相当于把几箱货搬到同一个货架上,能搬就搬,不强制重新分拣。在执行层面,它只是改变了 RDD 分区与子分区的映射关系,某些子分区可能对应多个父分区,但父分区与子分区的 blood relation 是一对多的窄依赖,不需要跨节点传输数据。
但coalesce的“快”是有前提的。当合并比例特别悬殊时——比如 5000 个分区合并到 10 个——它并不会均匀地把 500 个父分区分配给每个子分区,而是按照默认的 HashPartitioner 方式把连续的父分区合并到同一个子分区。这就会导致严重的倾斜问题:一个子分区可能对应 3000 个父分区,另一个只对应 200 个。数据分布不均,部分 task 执行时间远超平均线,整个 stage 被拖慢。
2.2 数据搬家的代价:repartition 为什么会 shuffle
repartition实际上调用了coalesce(n, shuffle = true),本质是强制把所有分区数据打散、重分区。它会通过RoundRobinPartitioning或你指定的Partitioner,把数据均匀地重新分配给 n 个分区。这个过程会触发真正的宽依赖(wide dependency),即所有数据节点之间需要跨网络传输,在 Spark UI 上表现为一个独立的 Shuffle Stage。
repartition的特点是均匀:因为每个父分区的数据都被打散,最后每个子分区的数据量会相对均衡,特别适合解决数据倾斜、key 分布不均、以及需要按特定字段重新组织数据(如 hash join 优化)的场景。代价也很直白:shuffle 一定要落盘(或走网络),磁盘 I/O、网络 I/O、序列化和反序列化的开销一并算上。
我在做网约车订单数据清洗时遇到过这样的场景:订单表按订单 ID 做了 hash 分区,但后续需要按司机 ID 做 join,此时分区方式与 join key 不匹配,导致大量数据倾斜。这时候repartition(col("driver_id"))就是合理的——重新按司机 ID 分桶,能让 join 阶段的数据分布更均衡。
2.3 两者区别的速查对照表
很多场景下,coalesce和repartition的选择并不难,关键是清楚自己的需求。我把两者的核心差异整理成一张对照表,平时写代码前扫一眼,基本不会选错:
| 维度 | coalesce | repartition |
|---|---|---|
| 源码本质 | CoalesceExec,窄依赖 | 调用coalesce(shuffle = true),宽依赖 |
| 是否触发 shuffle | 否 | 是 |
| 数据分布 | 不均匀(合并比例悬殊时倾斜明显) | 均匀(RoundRobin 或 Hash 重新分配) |
| 性能 | 快,开销小 | 慢,有额外的磁盘/网络 I/O |
| 适用方向 | 只减少分区数量 | 增加分区数量、解决倾斜、按 key 重分区 |
| 是否支持自定义分区字段 | 否 | 支持repartition(partitions, col) |
也别忽视spark.sql.shuffle.partitions。这个参数控制所有 shuffle 操作生成的分区数,很多人忽略了它对文件数量的影响。比如你写了df.groupBy("city").count()且没有动态分区的情况下,shuffle 默认生成 200 个分区,写出去就是 200 个文件。如果你想控制输出文件数量,必须在这个层面对齐:要么在聚合前把spark.sql.shuffle.partitions调小,要么聚合后合并分区。这一点后面在案例里会详细演示。
3. 实战选型:不同场景下合并参数怎么选
3.1 场景一:数据清洗后的小文件合并
网约车/农产品这类实时采集业务的源表普遍有大量小文件。任务流程通常是:读取源表 → 清洗(过滤脏数据、格式转换)→ 写回数仓。这类任务的瓶颈不是计算能力,而是下游查询时小文件过多导致的元数据开销。
我建议的处理方式是:清洗阶段用coalesce,在写完之前把分区数降到目标值。因为清洗阶段只涉及 map 类操作,没有发生宽依赖,用coalesce不会引入额外 shuffle,只需注意数据分布是否均匀。实操中,如果源表分区数 2000,你希望落地文件控制在 50 个以内,先filter后再coalesce(50)写入。但有个前提:清洗阶段如果有 filter,过滤后的数据量已经明显减少,此时合并比例不至于太悬殊,coalesce的倾斜问题可以被接受。
3.2 场景二:倾斜严重,需要增加分区或按 key 重分布
当某个热门城市的订单量是冷门城市的几百倍,按城市聚合时必然倾斜。此时coalesce解决不了问题,反而可能因为合并造成更严重的倾斜。正确做法是repartition(新分区数, 热点列或加盐列),甚至repartition(col("order_city"))。增加分区数的唯一方案是 repartition,因为coalesce不能增大分区(增大时不生效或退化为 shuffle)。
这里有个实践细节:如果热点键只有一个,单纯按 key repartition 依然会倾斜,因为 key 相同的记录都在同一分区。此时要考虑“加盐”策略:把热点 key 拆成多个子 key,处理后去盐。这个场景属于倾斜治理,本质已经超出了“合并参数区分”的范畴,但你要知道 repartition 是按 key 哈希的,相同 key 永远分布到同一个分区。
3.3 场景三:join 后的分区控制
经验教训:不要先 join 再大量合并,尽量在 join 前就把分区对齐。做过网约车项目数据开发的都有体会,订单表和司机表 join 后,如果两表的分区数不一致,shuffle 就会发生在 join 阶段,此时你再coalesce(20),只是第二次 shuffle 时的窄依赖优化,无法挽回第一次 shuffle 的巨大开销。
我的推荐方案:join 前先repartition对齐 join key 的分区数和分区方式,join 之后若仍需减少文件数,再用coalesce。比如:
val left = df1.repartition(100, col("driver_id")) val right = df2.repartition(100, col("driver_id")) val joined = left.join(right, Seq("driver_id"), "inner") val output = joined.coalesce(20)这个过程中,repartition(100, col("driver_id"))会产生第一个 shuffle stage,coalesce(20)不会产生额外 shuffle,但均匀性较差。如果你对均匀性有要求,output.repartition(20)会在已经按键对齐的基础上再 shuffle 一次,通常是没必要的奢侈。
3.4 场景四:写 Hive 表的动态分区
逻辑分区的坑不走一遍真的想不到。一个常见的坑是:动态分区表会根据分区列的值动态创建目录,此时你控制的分区数会被动态分区目录数量覆盖。比如你按日期分区写入 HDFS,最终文件数等于“写入时分区数 × 日期分区数”。
遇到这个场景,我的经验是:在写出之前不要过度合并,给动态分区预留足够的空间。常见做法是coalesce(每个动态分区的目标文件数 × 动态分区数)。但如果动态分区数本身很大(比如 500 个日期分区,你想每个分区 10 个文件,那就是 5000),此时宁可直接把spark.sql.shuffle.partitions调整为 5000,让每次 shuffle 后的输出自然对齐——不要用 coalesce 强行对到 5000(它可能产生严重倾斜)。这是参数调优与分区策略的交叉点,很多人忽略它的原因是只盯着“合并”这一步,而忘了下游的动态分区拆分。
4. 完整实操案例:从定位问题到参数落地
4.1 一个真实任务的现场回放
我最近维护的一个农产品价格预警任务就是这样。每天批量跑一次数据清洗和聚合,源表是十几个外部数据源导入的 Hive 表,每张表有一堆几十 KB 的小文件,总文件数约 6000。清洗后要按品类聚合,输出价格趋势表。任务初期跑完需要 40 分钟,其中近一半时间消耗在最后写表和元数据操作上。打开 Spark UI,清晰的瓶颈是最后的 save 阶段,写了 6000 个文件,每个文件平均只有 80KB。
经验丰富的开发一眼就明白:这是典型的分区数失控。但解决方案不能是一刀切的coalesce(50),那样清洗完、聚合完之后数据量可能只剩原来的十分之一,合并比例从 6000 到 50,coalesce必然倾斜。我需要先分析 stage 信息:哪些 stage 是 shuffle 产生的,哪些是 map 产生的。
4.2 通过 Spark UI 定位分区失控点
Spark UI 的 Stages 页签会清晰展示每个 stage 的 shuffle 读写量和输出分区数。我按以下步骤排查:
- 看 Event Timeline,找出耗时最长的 stage;
- 点进该 stage 的 Summary Metrics,查看分区数、shuffle read 量、各个 task 的执行时间分布;
- 对比 shuffle write 与实际分区大小,判断是否倾斜。
实测结果:漫长的耗时集中在聚合后的 save 阶段,shuffle write 总量约 1.2GB,分区数却有 6000 个。而聚合操作的输入只有 6000 个小文件,但配置的spark.sql.shuffle.partitions还是默认的 200,聚合本身没有产生异常文件数,问题出在源表的分区数直接被继承到了输出。注意,groupBy之后的 shuffle 默认是 200 个分区,这一步是好的,但后续如果有一个 map 操作(比如withColumn)会保留这 200 个分区,再写表时 Hive 的动态分区会按日期拆成 30 个分区,总共 6000 个输出文件。
所以我在代码里做了两件事:在聚合之前,先对源数据做一次 coalesce 到合理的分区数,比如清理后 6000 个文件合并到 300 个分区(每分区约 4MB,离理想大小有差距,但主要用于后续聚合,聚合阶段 shuffle 会再变一次);聚合完成后,用 coalesce(30) 直接控制最终输出文件数——因为聚合后的 1.2GB 数据分 30 个文件,每个文件约 40MB,可接受。
4.3 参数选择的落地细节
具体代码改造见下:
// 读取源表 val raw = spark.read.table("ods_market_price") // 清洗:过滤、格式转换 val cleaned = raw.filter(col("price").isNotNull && col("price") > 0) .withColumn("date_str", to_date(col("dt"))) // 清洗后合并分区,为聚合做准备 val merged = cleaned.coalesce(300) // 聚合 val aggregated = merged.groupBy("category", "date_str").agg(avg("price").as("avg_price")) // 最终输出,控制文件数量 val output = aggregated.coalesce(30) output.write.mode("overwrite").insertInto("dws_market_price_daily")为什么这里coalesce(300)放在聚合前?因为清洗阶段的 6000 个小文件直接进入聚合,每个分区数据量太小,聚合前的 shuffle 是 6000 个分区的数据同时参与,网络开销增大。合并到 300 后,shuffle 的输入分区变少,但数据量不变(仍是清洗后的全部数据),聚合阶段 shuffle 效率更高。280MB 左右每分区的数据量也接近 128~256MB 的理想区间。
coalesce(30)放在聚合后则是因为此时 1.2GB 的结果数据已经确定了分区数 200(来自spark.sql.shuffle.partitions),直接合并到 30 个分区,减少输出文件。
我特意不用repartition的原因是:清洗和聚合阶段没有 key 分布不均的问题,用coalesce就够了,额外 shuffle 纯属浪费。如果你在这个场景不放心,可以对比跑一次repartition(30),时间会明显多出 10%~20%,因为多了一次全量 shuffle。
4.4 调优后的效果对比
改造前,任务 40 分钟,输出 6000 个小文件,下游查询平均耗时 12 秒(文件元数据开销大)。改造后,任务 12 分钟,输出 30 个大文件(每个约 40MB),下游查询平均耗时 2 秒。提升了 3 倍多的任务效率,下游查询提升 6 倍。这个案例的核心不是“把合并参数调到多少”,而是先搞清楚分区在哪个阶段失控,再决定用哪种合并方式、在哪个阶段合并。
关于输出文件数目标,这里有个思考逻辑:文件数的多少取决于下游消费方式和数据量。如果下游是 Hive 表,每个文件 100~300MB 是合理区间;如果下游是 Spark 读,太多小文件会拖慢读取,太少则会降低并行度。
5. 常见问题与排查技巧实录
5.1 我用了 coalesce,为什么数据依旧倾斜?
最常见的原因有两个:合并比例过于悬殊,连续性合并导致分布不均;或者coalesce后仍保留了原先的哈希分区模式,某些 key 天然集中在特定分区。排查方式很简单:coalesce 后打印出每个分区的数据量,或者看 Spark UI 里相关 task 的 input size 与 duration 分布。如果明显不均衡,就说明要用repartition而不是coalesce——宁可多付一次 shuffle,也不要让单个任务拖垮整个 stage。
还有一点:如果你的父 RDD 已经是哈希分区的,coalesce 只会把多个父分区合并成一个子分区,并不能实现哈希的重打散。之前有篇文章提到一个案例:repartition(100)后coalesce(10),依然是均匀的,因为 repartition 的结果按 RoundRobin 分桶,coalesce 连续性合并后相对均衡;但如果直接把非均匀分区 coalesce,除非参数刚好合并比例接近整除,否则倾斜。
5.2 repartition 之后 shuffle 数据量暴增怎么办
repartition是全量 shuffle,数据量有可能比你预期的多。原因在于 shuffle read 会把所有节点的数据通过网络拉取,如果之前的 task 有溢写(spill),shuffle write 可能比内存中的数据大数倍。
我的经验是四步排查:先看 Spark UI 里 shuffle write 大小,确认是否异常;再看 spill 指标,内存不足导致溢写会增加 I/O;然后考虑增加 Executor 内存或调大spark.sql.shuffle.partitions,减小单个 shuffle 分区的数据量;最后,如果 repartition 后的下游不需要立即 shuffle,尽量把它与下一个算子合并,避免连续两次 shuffle。
举个例子,df.repartition(200).groupBy("key")就是一次多余的 shuffle,因为 groupBy 本身会根据 key 再次打散数据。这种连续 shuffle 是最隐蔽的性能杀手,代码简单但执行往往卡顿。
5.3 为什么 coalesce(1) 之后依然生成了多个文件
这个问题高频出现,尤其写 Hive 动态分区表时。coalesce(1)只能保证一个 Spark 分区写一个文件,但如果写入的是动态分区表,一个分区内如果包含多个不同的动态分区值,就会产生多个文件。比如 1 个 Spark 分区写到dt=2024-01-01和dt=2024-01-02两个目录,就是两个文件。
这种情况下,你需要重新梳理逻辑:要么把coalesce(1)放到最终写的步骤,确保一次数据落盘只有一个分区;要么用repartition(col("dt"))先按动态分区键分桶,再写出。动态分区的目录数和 Spark 分区数是两个维度,容易混淆。
5.4 千万别忽略 AQE 自动合并分区的影响
Spark 3.0 起的 AQE(Adaptive Query Execution)会在运行阶段自动合并 shuffle 后的中小分区,默认启用spark.sql.adaptive.enabled=true,其中spark.sql.adaptive.coalescePartitions.enabled=true会自动把平均数据量过小的分区合并到spark.sql.adaptive.advisoryPartitionSizeInBytes(默认 64MB)。
这就导致一个有趣的现象:你手动敲了repartition(200),但执行时 AQE 可能把你合并到几十个分区,你的预期失效了。如果你需要精确控制分区数(比如为了保证下游 Hive 表的文件数量),建议设置spark.sql.adaptive.coalescePartitions.enabled=false,或者在合并参数上设置repartition(200).coalesce(30)这种“先均匀、后收紧”的组合,让最终分区数的偏差控制在可接受范围。
我实际测试过:在某个网约车订单日活报表的 pipeline 中,AQE 自动合并后的分区数经常比期望少 20%~40%,文件大小则相应变大。如果你对文件大小的宽容度高,这其实是件好事;但如果你要继续做分桶表、并且后续有严格的分区数依赖,就必须显式关闭 AQE 的部分合并且自己管控分区。
5.5 控制文件数的额外技巧:结合 bucketing 和分区裁剪
最后补充一个小技巧。某些场景下,与其纠结 coalesce 和 repartition,不如直接用 bucketing 和 partitionBy 组合。写 Hive 表时bucketBy(n, key)能保证相同 key 落在相同文件,配合sortBy可以极大优化下游 join。这属于“在源头设计合理分区”的思路,比事后合并更高效。
比如农产品价格表按category分桶,每桶 4 个文件,后续 join 品类维表时就能直接走 bucket pruning,不用 shuffle 对齐。但这要求你在建表时就规划好分区策略,比起事后调参收到的收益更大。
5.6 一个常见的误区:分区数调越小越好
最后再强调一遍,不要陷入“越小越好”的陷阱。数据量 10GB、集群有 100 核,你coalesce(1)落盘,写入是快了,但下游查询时只有一个 task 在跑,并行度瞬间降为 0,读取速度比合并前还慢。合并的目的是匹配数据量与集群并行度的平衡点,而不是盲目追求文件数量少。
我在做数据服务接口项目时就踩过:一个 500GB 的表用coalesce(5)输出,每文件 100GB,下游 Spark 读它时产生了严重的网络抖落和 OOM。最终调整为coalesce(50),每文件 10GB,集群并行度上去了,任务反而快了一倍。换算逻辑很简单:目标文件数 ≈ 数据总量(GB)/ 每个目标文件合理大小(GB),再结合你的 Executor 总数微调——不是拍脑袋定的。
按我个人经验,如果你需要把 1.2GB、6000 个分区的结果写 Hive,目标是 30~50 个文件,coalesce(30)通常够用,但前提是合并比例别超过 100 倍太多,否则用 repartition 重新均匀化更稳。分区的管理永远是“动态平衡”而不是“一劳永逸”,每次任务都值得多花 10 秒钟想清楚:我到底在减少分区、还是在重新分布数据?想清楚这一点,合并参数的坑基本就避开了大半。