PySpark聚合结果仅6行,保存至HDFS却生成200+空文件求助
解决Spark输出HDFS时生成大量空文件的问题
你遇到的是Spark处理小数据集时的典型问题——明明结果只有6行数据,却生成了200多个空文件,核心原因是你的DataFrame分区数远大于实际数据行数,大部分分区没有数据,Spark依然会为每个分区生成一个输出文件。
问题根源
你的aggregate DataFrame在聚合过滤前,大概率保留了之前操作(比如从HDFS读取大文件)的默认分区数(200+),而过滤后的results只有6行数据,这些数据只占了极少数分区,剩下的空分区就对应了那些空文件。
快速解决方案
直接合并分区,让输出文件数和数据量匹配:
- 先确认当前分区数(可选,用于验证问题):
print(results.rdd.getNumPartitions())
运行后应该会输出200左右的数字,印证我们的判断。
- 合并分区后再保存:
用coalesce()方法(窄依赖,不会触发shuffle,效率更高)将分区数减少到1个(或者和数据行数一致的数量,比如6个,但1个文件更便于管理):
results = aggregate.filter(aggregate["count"] > 2500) # 合并到1个分区后再保存 results.coalesce(1).write.save("hdfs://your-target-path")
如果需要指定输出格式(比如CSV/Parquet),可以补充格式参数:
results.coalesce(1).write.format("csv") \ .option("header", "true") \ .save("hdfs://your-target-path")
额外说明
- 如果你的场景后续需要并行处理这个结果文件,可以把
coalesce(1)改成coalesce(6),让每一行数据对应一个分区(不过6行数据并行意义不大)。 - 也可以用
repartition(1)替代coalesce(1),但repartition会触发shuffle操作,对于小数据来说性能差异可以忽略,但coalesce更高效。
内容的提问来源于stack exchange,提问作者zipline86
相关产品推荐
相关产品推荐

