Spark中repartition操作为何比partitionBy写入速度更快?
1 认知偏差修正
你之前混淆了两个完全不同的partitionBy:
- 你印象中「仅本地排序、不需要shuffle」的是RDD API的
partitionBy()转换算子,它和DataFrame的repartition类似,属于数据转换操作,会触发shuffle将指定列相同值的数据归集到同一个内存分区。 - 你代码中用到的是
DataFrameWriter的partitionBy()写入参数,作用是写入时按指定列的值将数据输出到不同的子目录下,本身不会触发shuffle,但会带来极高的IO开销。
2 两种方案的执行逻辑差异
repartition方案(更快的原因)
你先调用repartition("partition")触发一次可控的shuffle,直接将相同partition值的所有数据归集到同一个Spark内存分区中。shuffle完成后,每个写入Task对应一个内存分区,这个分区内所有数据的partition值完全相同,因此每个Task只需要打开1个文件句柄写数据即可。100个目标分区最多只会生成100个输出文件,文件IO、压缩流开销都极低,哪怕shuffle有一定开销,也远低于后续写入的开销节省。
写入端partitionBy方案(更慢的原因)
该方案没有提前做数据shuffle,你读入9万个XML文件时,Spark默认会生成数千甚至上万个输入分区,每个输入分区中大概率包含全部100种partition值的数据。此时每个处理输入分区的Task,都需要为每一种出现的partition值单独打开一个带gzip压缩的文件流写数据,如果你有1000个输入Task,同时打开的文件流就高达10万级,会带来极高的内存、IO开销,最终还会生成10万个级别的小文件,存储系统处理小文件的额外开销进一步拉长了执行时间。这种场景下,写入的开销已经远远超过一次shuffle的开销,因此速度远慢于先repartition再写入的方案。
3 适用场景区分
写入端partitionBy并不是所有场景都慢:如果你的目标分区值很少(比如按天分区,每个输入分区仅包含12天的数据),每个Task只需要打开12个文件流,此时不需要shuffle的partitionBy确实会更快。但如果目标分区数较多、且每个输入分区都覆盖大多数目标分区值,先repartition再写入的方案效率会高很多。
内容的提问来源于stack exchange,提问作者Robin Zimmerman

