PySpark使用repartition/partitionBy无法按列值输出单文件问题求解
按分区值输出唯一文件的最优实现方案
核心问题根源
你遇到的每个category分区文件夹下生成多个文件的问题,核心原因是Spark默认哈希分区规则:当repartition仅指定分区列未指定总分区数时,会使用spark.sql.shuffle.partitions的默认值200作为总分区数,如果Spark版本低于2.4,会存在同分区列值散列到多个分区的bug,最终导致写操作时同一个分区值对应多个输出文件。
最快方案(Spark 2.4+ 推荐,无额外Job、无多余开销)
不需要提前统计分区列的唯一值数量,仅需将spark.sql.shuffle.partitions参数设置为略大于业务场景中该列的最大可能唯一值数即可,无需精确匹配。额外的空分区会在写阶段自动跳过,不会生成空文件,也不会带来额外计算成本。
示例代码:
# 配置值只要≥你实际的category唯一值数量即可,比如你业务中最多有500个分类就设为500 spark.conf.set("spark.sql.shuffle.partitions", 500) # 仅按category列重分区,无需指定分区数 df.repartition("category").write.partitionBy("category").save("你的输出路径")
零Shuffle最优方案(仅适用于上游DataFrame已按category分区的场景)
如果你的上游计算逻辑已经是按category做分区处理,不需要额外重分区,直接在写操作时配置maxRecordsPerFile为大于单个category最大记录数的数值即可,完全避免Shuffle,是所有场景下速度最快的方案:
df.write.option("maxRecordsPerFile", 1000000000).partitionBy("category").save("你的输出路径")
Spark 2.x 兼容低开销方案
如果你使用的是Spark 2.3及以下版本,可以用采样预估的方式替代全量count distinct,统计开销比全量统计低90%以上:
# 采样10%数据预估唯一值数量,乘以1.2冗余系数避免预估偏小 estimated_distinct_cnt = df.select("category").sample(0.1).distinct().count() * 1.2 df.repartition(int(estimated_distinct_cnt), "category").write.partitionBy("category").save("你的输出路径")
方案避坑说明
- 不要使用
repartition(1, 'category')的写法:该写法会强制所有数据shuffle到单个节点处理,完全丧失并行度,大数据量下极易出现单节点内存溢出,执行效率极低。 - 不要使用全量
count distinct获取分区数的写法:该写法会触发一次全表扫描的额外Job,若数据未缓存会带来极大的额外开销,完全没有必要。
内容的提问来源于stack exchange,提问作者Alain
相关产品推荐
相关产品推荐

