
“我把一个 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 为什么会 shufflerepartition实际上调用了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 stagecoalesce(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 到 50coalesce必然倾斜。我需要先分析 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 有溢写spillshuffle 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 分区写到dt2024-01-01和dt2024-01-02两个目录就是两个文件。这种情况下你需要重新梳理逻辑要么把coalesce(1)放到最终写的步骤确保一次数据落盘只有一个分区要么用repartition(col(dt))先按动态分区键分桶再写出。动态分区的目录数和 Spark 分区数是两个维度容易混淆。5.4 千万别忽略 AQE 自动合并分区的影响Spark 3.0 起的 AQEAdaptive Query Execution会在运行阶段自动合并 shuffle 后的中小分区默认启用spark.sql.adaptive.enabledtrue其中spark.sql.adaptive.coalescePartitions.enabledtrue会自动把平均数据量过小的分区合并到spark.sql.adaptive.advisoryPartitionSizeInBytes默认 64MB。这就导致一个有趣的现象你手动敲了repartition(200)但执行时 AQE 可能把你合并到几十个分区你的预期失效了。如果你需要精确控制分区数比如为了保证下游 Hive 表的文件数量建议设置spark.sql.adaptive.coalescePartitions.enabledfalse或者在合并参数上设置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 秒钟想清楚我到底在减少分区、还是在重新分布数据想清楚这一点合并参数的坑基本就避开了大半。