news 2026/9/29 16:33:15

Spark合并分区全解析:coalesce与repartition的底层原理及实战选型

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Spark合并分区全解析:coalesce与repartition的底层原理及实战选型

“我把一个 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的选择并不难,关键是清楚自己的需求。我把两者的核心差异整理成一张对照表,平时写代码前扫一眼,基本不会选错:

维度coalescerepartition
源码本质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 读写量和输出分区数。我按以下步骤排查:

  1. 看 Event Timeline,找出耗时最长的 stage;
  2. 点进该 stage 的 Summary Metrics,查看分区数、shuffle read 量、各个 task 的执行时间分布;
  3. 对比 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 秒钟想清楚:我到底在减少分区、还是在重新分布数据?想清楚这一点,合并参数的坑基本就避开了大半。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 16:32:23

Paperclip:轻量级本地AI Agent协同协议实战

1. 项目概述:Paperclip 不是回形针,而是一个正在被误读的 AI 工程实践符号最近在掘金、知乎和 GitHub Trending 上频繁刷到“paperclip”这个词,点进去却发现不是 Office 文档里的那个金属小物件,也不是某款设计工具的代号&#x…

作者头像 李华
网站建设 2026/9/29 16:30:16

华为EC6108V9/V9C U盘卡刷教程:固件选择、操作步骤与救砖指南

说实话,在智能盒子圈子里,华为EC6108V9和EC6108V9C算是“骨灰级”的常青树了。2015年前后运营商大范围集采,让这款盒子走进了千万家庭,直到现在二手市场还经常能看到它的身影。但问题也随之而来:系统停留在Android 4.4…

作者头像 李华
网站建设 2026/9/29 16:29:44

OpenHarmony上Flutter跨平台实战:看书记录App列表模块

把吃灰很久的一块 OpenHarmony 开发板搬到桌上,插上电,打开 DevEco Studio 和终端的那一刻,我就知道这次不是玩票:我准备在上面跑一个看书记录 App。这个 App 的核心功能听上去很简单——把我手头同时在读的几本书管理起来&#x…

作者头像 李华
网站建设 2026/9/29 16:29:28

Wyse 3040瘦客户机固件升级与Horizon协议适配实战

1. 为什么100块的Wyse 3040不是“电子垃圾”,而是可复用的瘦客户机黄金备件你刷到这个标题时,第一反应可能是:“Dell Wyse 3040?那不是2016年就停产的老古董吗?100块买回来能干啥?当U盘挂件?”—…

作者头像 李华
网站建设 2026/9/29 16:26:49

Jetpack Compose SubcomposeLayout:测量驱动组合的自定义布局终极指南

在 Jetpack Compose 里做自定义布局,绝大多数场景一个 Layout 就够了:拿到子项,measure 一遍,然后按自己的规则摆放。但如果你遇到“要先知道文字实际占了多高,才决定要不要在右下角放一个展开按钮”“要根据每个标签…

作者头像 李华
网站建设 2026/9/29 16:26:37

数据依赖与数据库范式:从函数依赖到3NF/BCNF的规范化实战

有些场景其实不该无脑上3NF,这个我们后面再细说。1. 内容整体设计与思路拆解1.1 为什么先讲“数据依赖”而非直接讲“范式”我在带新人或者帮朋友团队做数据库评审的时候,发现一个很普遍的问题:很多人背得出三大范式(1NF、2NF、3N…

作者头像 李华