Spark转换数据后如何覆盖S3同路径同名CSV文件
问题根因
Spark默认overwrite写入模式是目录级操作:写入时先生成带随机UUID的临时staging目录,作业完成后会清空目标路径下所有原有内容,再将临时目录内Spark自动命名的part-*格式文件移入目标路径,既无法保留原文件名,也做不到按单个原文件逐一定向覆盖。
要实现按原同名文件逐个覆盖的需求,核心逻辑分两步:
- 读取数据时保留每条记录对应的源文件路径元信息
- 按源文件路径分组,每组数据先写入临时路径,再通过Hadoop S3文件系统API删除对应原文件,将临时输出的CSV重命名为原文件名完成替换,最后清理临时文件
PySpark 实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import input_file_name, col import os spark = SparkSession.builder.appName("S3CsvOverwrite").getOrCreate() hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() # 配置源路径 source_path = "s3://MyBucket/object/" # 1. 读取CSV时携带每条记录所属的源文件路径 df = spark.read.csv( f"{source_path}/*.csv", header=True, inferSchema=True ).withColumn("source_file", input_file_name()) # 2. 在此处执行你的自定义数据转换逻辑 transformed_df = df # 3. 逐文件处理覆盖 source_files = transformed_df.select("source_file").distinct().toPandas()["source_file"].tolist() temp_suffix = os.urandom(4).hex() temp_root = f"{source_path}_temp_overwrite_{temp_suffix}/" for file_path in source_files: file_name = os.path.basename(file_path) # 过滤当前文件对应的数据,移除元列 single_file_df = transformed_df.filter(col("source_file") == file_path).drop("source_file") # 单分区写入临时路径,保证输出单个CSV文件 single_file_df.coalesce(1).write.mode("overwrite").option("header", True).csv(f"{temp_root}{file_name}") # 操作S3文件系统完成替换 fs_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(file_path) fs = fs_path.getFileSystem(hadoop_conf) temp_dir_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(f"{temp_root}{file_name}") # 定位临时目录下生成的CSV结果文件 for file_status in fs.listStatus(temp_dir_path): temp_file_name = file_status.getPath().getName() if temp_file_name.startswith("part-") and temp_file_name.endswith(".csv"): fs.delete(fs_path, False) fs.rename(file_status.getPath(), fs_path) break # 清理全部临时文件 temp_root_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(temp_root) fs.delete(temp_root_path, True)
Scala 实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.input_file_name import org.apache.hadoop.fs.Path import java.security.SecureRandom object S3CsvOverwrite { def main(args: Array[String]): Unit = { val spark = SparkSession.builder().appName("S3CsvOverwrite").getOrCreate() val hadoopConf = spark.sparkContext.hadoopConfiguration val sourcePath = "s3://MyBucket/object/" val tempSuffix = SecureRandom.getInstanceStrong.nextLong().toHexString val tempRoot = s"${sourcePath}_temp_overwrite_$tempSuffix/" // 1. 读取CSV时携带每条记录所属的源文件路径 val df = spark.read .option("header", "true") .option("inferSchema", "true") .csv(s"$sourcePath/*.csv") .withColumn("source_file", input_file_name()) // 2. 在此处执行你的自定义数据转换逻辑 val transformedDf = df // 3. 逐文件处理覆盖 import spark.implicits._ val sourceFiles = transformedDf.select("source_file").distinct().as[String].collect() sourceFiles.foreach(filePath => { val fileName = filePath.split("/").last val singleFileDf = transformedDf.filter($"source_file" === filePath).drop("source_file") // 单分区写入临时路径,保证输出单个CSV文件 singleFileDf.coalesce(1).write.mode("overwrite") .option("header", "true") .csv(s"$tempRoot$fileName") // 操作S3文件系统完成替换 val fsPath = new Path(filePath) val fs = fsPath.getFileSystem(hadoopConf) val tempDirPath = new Path(s"$tempRoot$fileName") val tempFiles = fs.listStatus(tempDirPath) tempFiles.foreach(fileStatus => { val tempFileName = fileStatus.getPath.getName if (tempFileName.startsWith("part-") && tempFileName.endsWith(".csv")) { fs.delete(fsPath, false) fs.rename(fileStatus.getPath, fsPath) } }) }) // 清理全部临时文件 val tempRootPath = new Path(tempRoot) val fs = tempRootPath.getFileSystem(hadoopConf) fs.delete(tempRootPath, true) } }
注意事项
- 代码中
coalesce(1)用于保证每个源文件对应输出单个CSV,和原文件结构一致;如果单文件数据量极大出现OOM,可以调整分区数后自行合并输出文件,核心替换逻辑不变 - 需保证Spark任务绑定的IAM权限/密钥对包含对应S3路径的
s3:ListBucket、s3:PutObject、s3:DeleteObject权限,否则删除、重命名操作会报错 - 如果你的CSV有自定义分隔符、编码、跳过行等特殊配置,读写时需补全对应option参数,保证读写格式一致
内容的提问来源于stack exchange,提问作者Priya p
相关产品推荐
相关产品推荐

