如何在PySpark中写入CSV文件时不创建文件夹且指定自定义文件名
PySpark生成指定名称单CSV文件的实现方法
Spark作为分布式计算框架,默认多分区并行写入的特性决定了直接写入会生成文件夹和多分区随机名称文件,以下是仅用PySpark实现单指定名称CSV的方案,无需调用pandas:
实现逻辑
- 先将DataFrame合并为单分区,保证写入临时路径时仅生成1个CSV数据文件,优先用
coalesce(1)避免不必要的shuffle开销 - 调用Spark原生集成的Hadoop文件系统API,将临时文件夹内的唯一CSV文件移动到目标路径并重命名,最后清理临时文件夹
完整代码示例
from pyspark.sql import SparkSession from py4j.java_gateway import java_import # 初始化SparkSession spark = SparkSession.builder.appName("SingleCSVOutput").getOrCreate() # 替换为你自己业务的DataFrame df = spark.createDataFrame([(1, "张三"), (2, "李四"), (3, "王五")], ["id", "name"]) # 临时输出路径,不要和最终目标路径重名 temp_output_path = "/tmp/temp_csv_dir" # 最终自定义名称的CSV文件路径 target_csv_path = "/tmp/custom_name.csv" # 单分区写入临时路径 df.coalesce(1).write \ .option("header", "true") \ # 可根据需求调整是否写入表头、编码等参数 .mode("overwrite") \ .csv(temp_output_path) # 调用Hadoop API处理文件移动重命名 java_import(spark._jvm, 'org.apache.hadoop.fs.Path') hadoop_conf = spark._jsc.hadoopConfiguration() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) # 匹配临时目录内的CSV文件 temp_file_list = fs.listStatus(spark._jvm.Path(temp_output_path)) source_csv_file = None for file in temp_file_list: file_name = file.getPath().getName() if file_name.endswith(".csv"): source_csv_file = file.getPath() break if source_csv_file: # 移动并重命名到目标路径 fs.rename(source_csv_file, spark._jvm.Path(target_csv_path)) # 删除临时文件夹 fs.delete(spark._jvm.Path(temp_output_path), True) else: raise Exception("临时输出目录未找到生成的CSV文件")
注意事项
- 该方案仅适合数据量不大的场景,合并为单分区会将所有数据集中到单个Executor处理,数据量过大会触发内存溢出
- 写入对象存储(S3、OSS等)场景同样适用,只要Spark配置了对应存储的访问权限即可
- 如果DataFrame本身就是单分区,可省略
coalesce(1)步骤
内容的提问来源于stack exchange,提问作者cookie
相关产品推荐
相关产品推荐

