You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 13:45:06