You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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}")
    }
  }
}

关键步骤说明

  1. 递归遍历文件:通过自定义listAllFiles方法,自动递归进入每个子文件夹,同时过滤掉不需要的_SUCCESS标记文件。
  2. 提取分区信息:从文件路径中解析出DataPartition=后面的分区值,这一步依赖于你的路径结构,如果结构有变化,只需调整分割或正则逻辑即可。
  3. 生成时间戳:使用SimpleDateFormat生成符合要求的时间戳格式,确保每次运行生成的文件名唯一。
  4. 文件移动与重命名:利用Hadoop FS的rename方法完成文件迁移,S3环境下这个操作本质是跨对象存储的复制+删除,确保你的账号有足够的读写权限。

注意事项

  • 如果原文件是GZIP压缩格式,而目标文件需要是未压缩的TXT,你需要在移动前添加解压逻辑(可以用Spark读取原文件再写入,或者用Hadoop的压缩API处理)。
  • 确保Spark的Hadoop配置正确配置了S3的访问密钥(如果是非公开桶),避免出现权限不足的报错。
  • 若待处理文件数量极大,建议添加批量处理逻辑,避免单线程遍历带来的性能瓶颈。

内容的提问来源于stack exchange,提问作者Atharv Thakur

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:50:39