Spark使用repartition与coalesce输出DataFrame为CSV的问题咨询
Spark单分区排序输出单CSV异常问题原理解析
核心背景规则
Spark SQL中触发shuffle的算子(包括orderBy、groupBy、join等),shuffle后的默认分区数由配置项spark.sql.shuffle.partitions控制,默认值为200,该规则优先级高于输入DataFrame的原有分区数。
各异常现象底层逻辑
- 先
orderBy再调用repartition输出仅分块有序:orderBy本身会生成全局有序的多分区结果(分区之间有序、分区内部也有序),但后续的repartition会将已排序的数据再次打散重分区,仅能保证新的每个分区内部有序,跨分区的全局顺序被破坏,最终结果就是分块有序。 - 移除
repartition后输出200个CSV且仅能查到部分记录:哪怕输入DataFrame只有1个分区,orderBy触发shuffle时,默认会生成200个分区存储排序结果,输出时每个分区对应一个CSV文件,200个文件的总数据才是完整结果,仅读取单个文件自然只能看到部分记录。 orderBy前调用各类repartition仍输出200个文件:repartition是shuffle算子,在orderBy前调用仅能修改shuffle前的map端分区数,orderBy触发自身的shuffle时,仍然会按照默认200个分区生成结果,所以最终输出还是200个文件。orderBy前调用coalesce(1)问题修复:coalesce(1)属于窄依赖,不会触发shuffle,会直接将原数据合并为单个分区,后续调用orderBy时,单分区内排序不需要触发shuffle,不会用到默认200分区的规则,排序后仍然是单分区,输出自然就是单个全局有序的CSV文件。
coalesce与repartition的核心差异
repartition(numPartitions):无论调大还是调小分区数,都会触发shuffle,将数据按哈希规则打散分配到指定数量的分区,适合需要均匀打散数据的场景。coalesce(numPartitions):仅在调小分区数时生效,属于窄依赖,不会触发shuffle,仅合并相邻分区的数据、不会打乱原有数据顺序,开销远低于repartition,适合缩小分区数的场景。
全局有序输出单CSV的最优方案
根据数据量大小可以选择两种实现:
- 小数据量场景:先执行
coalesce(1)合并为单分区,再执行orderBy排序,无额外shuffle开销,输出即为单个有序CSV。 - 大数据量场景:先执行
orderBy并行排序(可提前调大spark.sql.shuffle.partitions提升排序并行度),再执行coalesce(1)合并分区,此时合并操作不会打乱orderBy生成的全局顺序,兼顾排序性能和输出要求。
内容的提问来源于stack exchange,提问作者Ed Cheng
相关产品推荐
相关产品推荐

