PySpark使用partitionBy写parquet时如何控制总输出文件数为指定值
解决方案
问题根因
你之前的写法df.repartition(1200).write.partitionBy("mykey").parquet(...)会生成大量小文件的原因是:repartition(1200)是将全量数据随机打散为1200个Spark分区,每个Spark分区内都包含所有mykey取值的部分数据,写入时每个Spark分区会为每个mykey分区单独生成一个文件,最终总文件数为 1200 * mykey唯一值数量,和你的预期不符。
最优兼容方案(Spark 2.4 完全支持)
你可以通过哈希分片的方式提前控制数据分布,保证总文件数不超过1200,同时避免OOM风险,代码示例如下:
from pyspark.sql.functions import col, hash EXPECTED_TOTAL_FILES = 1200 df.withColumn("shard_key", hash(col("mykey")) % EXPECTED_TOTAL_FILES) \ .repartition(EXPECTED_TOTAL_FILES, "shard_key") \ .drop("shard_key") \ .write \ .partitionBy("mykey") \ .parquet("你的输出路径")
方案说明
- 该逻辑通过对
mykey做哈希取模,将同一个mykey的所有数据路由到同一个或少数几个Spark分区中,写入时每个Spark分区最多只会给每个mykey生成1个文件,最终总文件数严格控制在1200以内,完全匹配你的需求 - 不需要提前统计
mykey的唯一值数量,比你想到的按唯一值数量做除法的方案更优雅 - Spark分区数始终保持为1200,和你最初的配置一致,不会出现单分区数据量过大导致的内存溢出问题
- 所有用到的函数和API均为Spark 2.4原生支持,不需要升级版本
替代备选方案
如果你的场景允许事后处理,也可以在写入完成后调用Hadoop FileSystem的API对同个mykey下的小文件做合并,不过该方案会产生额外的IO开销,更推荐使用上述提前分片的方案。
内容的提问来源于stack exchange,提问作者safex
相关产品推荐
相关产品推荐

