如何将Spark DataFrame以指定名称的单个CSV文件写入S3(无文件夹)
Spark DataFrame写入S3生成指定名称单个CSV的解决方案
Spark默认CSV写入机制会生成带随机名称的part文件及元数据文件(如_SUCCESS、.crc),即使使用coalesce(1)合并分区也无法直接指定输出文件名。以下是Java/Scala环境下的解决方法:
核心思路
先将DataFrame写入S3临时目录,找到生成的单个part文件后重命名为指定名称,最后清理临时目录及冗余文件。
Scala 实现代码
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.spark.sql.SaveMode // 定义目标文件路径和临时目录 val targetCsvPath = "s3a://<your-bucket>/<target-folder>/custom-name.csv" val tempDirPath = "s3a://<your-bucket>/<target-folder>/temp-" + System.currentTimeMillis() // 1. 将数据写入临时目录 dataFrame.coalesce(1) .write .mode(SaveMode.Overwrite) .option("header", "true") .csv(tempDirPath) // 2. 获取Hadoop文件系统实例 val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration) // 3. 定位临时目录下的part文件 val partFile = fs.listStatus(new Path(tempDirPath)) .filter(_.getPath.getName.startsWith("part-")) .map(_.getPath) .headOption // 4. 重命名part文件到目标路径并清理临时目录 partFile.foreach { file => fs.rename(file, new Path(targetCsvPath)) fs.delete(new Path(tempDirPath), true) // 递归删除临时目录 }
Java 实现代码
import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SaveMode; // 定义目标文件路径和临时目录 String targetCsvPath = "s3a://<your-bucket>/<target-folder>/custom-name.csv"; String tempDirPath = "s3a://<your-bucket>/<target-folder>/temp-" + System.currentTimeMillis(); // 1. 将数据写入临时目录 dataFrame.coalesce(1) .write() .mode(SaveMode.Overwrite) .option("header", "true") .csv(tempDirPath); // 2. 获取Hadoop文件系统实例 FileSystem fs = FileSystem.get(spark.sparkContext().hadoopConfiguration()); // 3. 定位临时目录下的part文件 Path tempDir = new Path(tempDirPath); Path partFile = null; for (var status : fs.listStatus(tempDir)) { Path path = status.getPath(); if (path.getName().startsWith("part-")) { partFile = path; break; } } // 4. 重命名文件并清理临时目录 if (partFile != null) { fs.rename(partFile, new Path(targetCsvPath)); fs.delete(tempDir, true); // 递归删除临时目录 }
注意事项
- 权限配置:确保Spark集群已配置S3读写权限,包含
hadoop-aws依赖及AWS凭证(如Access Key/Secret Key或IAM角色)。 - 数据量风险:
coalesce(1)会将所有数据集中到单个Executor节点,数据量过大时可能引发内存溢出,超大数据集建议先通过其他方式合并,再重命名。 - 临时目录唯一性:临时目录加入时间戳后缀,避免多任务执行时的目录冲突。
内容的提问来源于stack exchange,提问作者Ritik Kamra
相关产品推荐
相关产品推荐

