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

