如何将PySpark DataFrame保存为自定义文件名的CSV文件?
问题结论
可以实现自定义文件名,但无法直接通过DataFrameWriter.csv()的参数直接指定最终文件名,需要根据你的数据量和业务场景选择对应方案:
方案1:输出单文件 MyDataFrame.csv(仅适合小数据量场景)
Spark默认多分区并行写入会生成多个part-前缀的文件,如果你需要合并为单个自定义名称的文件,步骤如下:
- 先将DataFrame合并为单分区,写入临时目录
- 调用Hadoop FileSystem API将临时目录下生成的唯一part文件重命名为目标文件名,移动到指定路径后清理临时目录
PySpark实现代码
# 定义临时路径与最终目标文件路径 tmp_save_dir = "/your/tmp/dir" target_file_path = "/your/target/path/MyDataFrame.csv" # 单分区写入临时目录 MyDataFrame.coalesce(1).write.csv( tmp_save_dir, mode="overwrite", header="true" ) # 获取Hadoop FileSystem实例 sc = MyDataFrame.sql_ctx.sparkContext hadoop_conf = sc._jsc.hadoopConfiguration() fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) # 找到临时目录下生成的csv文件 part_file = list(fs.globStatus(sc._jvm.org.apache.hadoop.fs.Path(f"{tmp_save_dir}/part-*.csv")))[0].getPath() # 重命名到目标路径 fs.rename(part_file, sc._jvm.org.apache.hadoop.fs.Path(target_file_path)) # 删除临时目录 fs.delete(sc._jvm.org.apache.hadoop.fs.Path(tmp_save_dir), True)
注意:
coalesce(1)会将所有数据拉到同一个Executor处理,数据量过大会导致内存溢出、写入效率极低,大数据量场景不建议使用。
方案2:多文件自定义前缀(适合大数据量场景)
如果不需要合并为单文件,仅需要替换默认的part-前缀为MyDataFrame,可以直接使用Spark 2.2及以上版本支持的outputFileName配置:
实现代码
MyDataFrame.write \ .option("outputFileName", "MyDataFrame.csv") \ .mode("overwrite") \ .csv(csv_path, header="true")
写入后生成的文件格式为 MyDataFrame-00000-<随机UUID>-c000.csv、MyDataFrame-00001-<随机UUID>-c000.csv,每个分区对应一个文件,不会影响分布式写入性能。
底层逻辑说明
Spark原生不支持直接指定单个输出文件名,是因为分布式架构下多个Executor会并行写入数据,多个进程同时写入同一个文件会出现资源冲突、数据覆盖问题,因此默认设计为每个分区写入独立文件,自定义文件名需要通过后处理重命名实现。
内容的提问来源于stack exchange,提问作者AlwaysWondering
相关产品推荐
相关产品推荐

