Spark Scala实现基于S3文件夹名的文件重命名与移动
解决Spark输出文件递归遍历S3子文件夹并重命名的问题
你的问题核心在于listStatus默认只会列出当前目录的直接内容,不会递归遍历子文件夹。下面我会给出完整的解决方案,包含递归遍历文件、提取分区信息、生成指定格式文件名以及移动文件的代码:
完整代码实现
import org.apache.hadoop.fs._ import java.text.SimpleDateFormat import java.util.Date // 初始化路径和配置 val srcRoot = new Path("s3://trfsmallfffile/Segments/output") val destRoot = new Path("s3://trfsmallfffile/Segments/Finaloutput") val conf = sc.hadoopConfiguration val fs = srcRoot.getFileSystem(conf) // 递归遍历所有子文件的工具方法 def listAllFiles(fs: FileSystem, path: Path): Array[Path] = { val statuses = fs.listStatus(path) statuses.flatMap { status => if (status.isDirectory) { // 如果是文件夹,递归遍历子目录 listAllFiles(fs, status.getPath) } else { // 如果是文件,过滤掉_SUCCESS文件 val fileName = status.getPath.getName if (fileName != "_SUCCESS") Array(status.getPath) else Array.empty[Path] } } } // 获取所有需要处理的文件 val allFiles = listAllFiles(fs, srcRoot) // 生成当前时间戳(格式:yyyy-MM-dd-HHmm) val timestamp = new SimpleDateFormat("yyyy-MM-dd-HHmm").format(new Date()) // 遍历处理每个文件 allFiles.foreach { srcPath => // 1. 提取DataPartition的值(比如从路径中解析出Japan/SelfSourcedPrivate等) val pathStr = srcPath.toString val partitionValue = pathStr.split("DataPartition=")(1).split("/")(0) // 2. 提取原文件的标识部分(示例中的1971-BAL.1,这里假设从原part文件名提取,可根据实际调整) // 你可以根据自己的实际文件名规则修改这部分逻辑 val originalFileId = srcPath.getName.split("\\.")(0).split("-")(1) // 3. 构建新文件名 val newFileName = s"Fundamental.FinancialStatement.FinancialStatementLineItems.${partitionValue}.${originalFileId}.${timestamp}.Full.txt" val destPath = new Path(destRoot, newFileName) // 4. 移动文件(S3上的rename本质是复制后删除原文件) if (fs.exists(destPath)) { println(s"目标文件${destPath}已存在,跳过") } else { val success = fs.rename(srcPath, destPath) if (success) { println(s"成功移动文件:${srcPath} -> ${destPath}") } else { println(s"移动文件失败:${srcPath}") } } }
关键步骤说明
- 递归遍历文件:通过自定义
listAllFiles方法,自动递归进入每个子文件夹,同时过滤掉不需要的_SUCCESS标记文件。 - 提取分区信息:从文件路径中解析出
DataPartition=后面的分区值,这一步依赖于你的路径结构,如果结构有变化,只需调整分割或正则逻辑即可。 - 生成时间戳:使用
SimpleDateFormat生成符合要求的时间戳格式,确保每次运行生成的文件名唯一。 - 文件移动与重命名:利用Hadoop FS的
rename方法完成文件迁移,S3环境下这个操作本质是跨对象存储的复制+删除,确保你的账号有足够的读写权限。
注意事项
- 如果原文件是GZIP压缩格式,而目标文件需要是未压缩的TXT,你需要在移动前添加解压逻辑(可以用Spark读取原文件再写入,或者用Hadoop的压缩API处理)。
- 确保Spark的Hadoop配置正确配置了S3的访问密钥(如果是非公开桶),避免出现权限不足的报错。
- 若待处理文件数量极大,建议添加批量处理逻辑,避免单线程遍历带来的性能瓶颈。
内容的提问来源于stack exchange,提问作者Atharv Thakur
相关产品推荐
相关产品推荐

