如何在Dataproc Serverless运行的Spark中重命名GCS文件?
问题描述
将Spark DataFrame写入文件后,使用以下Scala代码重命名文件:
val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration) val file = fs.globStatus(new Path(path + "/part*"))(0).getPath().getName() fs.rename(new Path(path + "/" + file), new Path(path + "/" + fileName))
本地运行Spark时代码正常,但在Dataproc上执行Jar包时抛出错误:
Exception in thread "main" java.lang.IllegalArgumentException: Wrong bucket: prj-***, in path: gs://prj-*****/part*, expected bucket: dataproc-temp-***
推测文件需等待任务结束后才会保存到目标桶,导致重命名失败。尝试配置.option("mapreduce.fileoutputcommitter.algorithm.version", "2")后问题仍未解决,进一步发现spark.sparkContext.hadoopConfiguration似乎强制要求基础桶为dataproc-temp-*类型,完整堆栈跟踪如下:
Exception in thread "main" java.lang.IllegalArgumentException: Wrong bucket: prj-**, in path: gs://p**, expected bucket: dataproc-temp-u*** at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem.checkPath(GoogleHadoopFileSystem.java:95) at org.apache.hadoop.fs.FileSystem.makeQualified(FileSystem.java:667) at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystemBase.makeQualified(GoogleHadoopFileSystemBase.java:394) at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem.getGcsPath(GoogleHadoopFileSystem.java:149) at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystemBase.globStatus(GoogleHadoopFileSystemBase.java:1085) at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystemBase.globStatus(GoogleHadoopFileSystemBase.java:1059)
解决方案
1. 针对目标路径创建独立FileSystem实例
Dataproc默认的spark.sparkContext.hadoopConfiguration绑定到临时桶(dataproc-temp-*),直接用它获取的FileSystem只能操作该临时桶。需针对目标GS路径单独创建FileSystem实例:
import org.apache.hadoop.fs.Path import org.apache.hadoop.fs.FileSystem import org.apache.hadoop.conf.Configuration // 基于现有配置创建独立配置实例 val conf = new Configuration(spark.sparkContext.hadoopConfiguration) val targetPath = new Path(path) // 获取目标路径对应的FileSystem val fs = targetPath.getFileSystem(conf) // 执行重命名逻辑 val partFile = fs.globStatus(new Path(path + "/part*"))(0).getPath() fs.rename(partFile, new Path(path + "/" + fileName))
2. 确保输出提交器配置生效
搭配上述修改,全局设置输出提交器版本,让文件直接写入目标路径,避免临时目录同步延迟:
// 配置提交器版本 spark.conf.set("mapreduce.fileoutputcommitter.algorithm.version", "2") // 针对Parquet格式指定直接提交器 spark.conf.set("spark.sql.parquet.output.committer.class", "org.apache.spark.sql.parquet.DirectParquetOutputCommitter")
3. 替代方案
- 小数据量场景:写入时直接合并为单个文件,无需后续重命名:
df.coalesce(1).write.format("parquet").save(path) - 任务后处理:任务结束后通过
gsutil命令重命名(可在Dataproc作业的后续步骤或初始化动作中执行):gsutil mv gs://prj-****/part-* gs://prj-****/your-target-filename.parquet
内容的提问来源于stack exchange,提问作者Daniel Fletemier
相关产品推荐
相关产品推荐

