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

Spark 2.4.0写入CSV至对象存储时跳过_temporary文件夹的方法

解决Spark 2.4.0写入CSV到对象存储时_temporary目录移动耗时过长的问题

问题背景

用PySpark 2.4.0将数据写入支持类Hadoop接口的对象存储,设置2048个分区生成小CSV文件,写入_temporary文件夹仅需6分钟,但把文件从临时目录移到最终目录要花1.5小时。试过以下方案都没用:

  • 设置spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=2:该算法仍会依赖临时文件
  • 修改Parquet专用的OutputCommitter:不适用于CSV格式
  • Zero Rename Committer:仅适配S3且可能需要更高版本Spark

核心原因

对象存储的文件重命名操作远慢于本地文件系统,2048个小文件的批量重命名会导致耗时剧增。Spark默认的FileOutputCommitter为保证写入原子性,会先把任务输出写到临时目录,等所有任务完成后再批量移动到最终目录,这套逻辑在对象存储场景下代价极高。

可行解决方案(Spark 2.4.0适配)

1. 自定义CSV OutputCommitter(最优方案)

Spark的CSV输出默认用FileOutputCommitter,我们可以自定义该类,让任务直接把输出写到最终目录,跳过全局临时目录的移动步骤。

步骤1:编写Java自定义Committer

因为Spark的OutputCommitter是Java接口,需要写Java类并打包成JAR:

import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.mapreduce.JobContext;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter;

import java.io.IOException;

public class DirectCSVOutputCommitter extends FileOutputCommitter {
    public DirectCSVOutputCommitter(Path outputPath, TaskAttemptContext context) throws IOException {
        super(outputPath, context);
    }

    @Override
    public void commitTask(TaskAttemptContext context) throws IOException {
        // 获取任务临时输出路径和最终目录路径
        Path taskTempPath = getWorkPath(context);
        Path finalOutputPath = getOutputPath();
        FileSystem fs = finalOutputPath.getFileSystem(context.getConfiguration());

        // 直接将任务输出文件移动到最终目录
        for (FileStatus status : fs.listStatus(taskTempPath)) {
            Path src = status.getPath();
            Path dest = new Path(finalOutputPath, src.getName());
            if (!fs.rename(src, dest)) {
                throw new IOException("无法重命名文件: " + src + " -> " + dest);
            }
        }
        // 删除任务临时目录
        fs.delete(taskTempPath, true);
    }

    @Override
    public void commitJob(JobContext context) throws IOException {
        // 清理全局临时目录(如果存在)
        Path tempDir = new Path(getOutputPath(), "_temporary");
        FileSystem fs = tempDir.getFileSystem(context.getConfiguration());
        if (fs.exists(tempDir)) {
            fs.delete(tempDir, true);
        }
    }
}

步骤2:打包JAR并在PySpark中引用

把上述Java代码编译打包成JAR文件(比如direct-csv-committer.jar),启动PySpark时通过--jars参数引入:

pyspark --jars /path/to/direct-csv-committer.jar

步骤3:配置Spark使用自定义Committer

在PySpark代码里设置CSV输出的Committer类:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("DirectCSVWrite") \
    .config("spark.sql.csv.output.committer.class", "com.yourpackage.DirectCSVOutputCommitter") \
    .getOrCreate()

# 先删除目标目录保证overwrite语义(跳过原子性检查后需手动保证)
dbutils = spark.sparkContext._jvm.com.databricks.dbutils.DBUtils(spark.sparkContext)
if dbutils.fs.exists(dsn):
    dbutils.fs.rm(dsn, recurse=True)

# 执行CSV写入
df.write.csv(
    dsn, 
    mode="overwrite", 
    sep=sep, 
    encoding=encoding, 
    quote="\"",
    escape="\""
)

注意:这个方案会牺牲全局原子性,如果任务中途失败,目标目录可能残留部分文件,所以需要提前手动删除目标目录来保证overwrite的语义。

2. 调整分区数量

减少分区数能降低需要重命名的文件总数,从而减少移动耗时。比如把分区数调到1024或512,在文件大小和数量之间找平衡,测试是否能显著缩短移动时间。

3. 本地写入后批量上传

如果集群节点有足够本地存储,可以先把CSV写到本地目录,再用对象存储的批量上传工具(比如hadoop fs -put或对应SDK)把文件批量传到目标目录,完全绕过Spark的临时目录机制。但这个方案需要处理分布式写入本地的协调问题,适合数据量较小的场景。


内容的提问来源于stack exchange,提问作者Timothy Wu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 22:45:46