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

如何用Spark/Scala(2.0版)将HDFS子目录文件复制至基础目录并归档子目录

Alright, let's tackle this problem step by step using Scala and Spark 2.x. I'll walk you through the core logic, code implementation, and key considerations to make this work smoothly.

核心任务拆解

First, let's break down exactly what we need to accomplish:

  • Traverse all date-named subdirectories under /user/srav/ (like 20190101, 20180101)
  • Copy all .dat files from these subdirectories up to the base /user/srav/ directory
  • Archive the original date subdirectories (either move them to an archive folder or compress them for cleanup)
Implementation with Scala & Spark 2.x

Since Spark runs on top of Hadoop, we can leverage Hadoop's FileSystem API to interact directly with HDFS—no extra dependencies required. Here's a complete, ready-to-use code example:

import org.apache.spark.sql.SparkSession
import org.apache.hadoop.fs.{FileSystem, Path, FileStatus}
import java.net.URI

object HdfsDatFileManager {
  def main(args: Array[String]): Unit = {
    // Initialize Spark Session (Spark 2.x standard setup)
    val spark = SparkSession.builder()
      .appName("HDFS Dat File Copier & Directory Archiver")
      .getOrCreate()
    
    // Get HDFS FileSystem instance (uses cluster's default FS config automatically)
    val hdfsUri = spark.sparkContext.hadoopConfiguration.get("fs.defaultFS")
    val fs = FileSystem.get(URI.create(hdfsUri), spark.sparkContext.hadoopConfiguration)
    
    // Define core paths
    val basePath = new Path("/user/srav/")
    val archiveRootPath = new Path("/user/srav/archive/")
    
    // Create archive root if it doesn't exist
    if (!fs.exists(archiveRootPath)) {
      fs.mkdirs(archiveRootPath)
      println(s"Created archive directory: ${archiveRootPath}")
    }

    // 1. Filter out valid date-named subdirectories
    val dateSubDirs = fs.listStatus(basePath)
      .filter(status => status.isDirectory && status.getPath.getName.matches("^\\d{8}$")) // Matches 8-digit date format
    
    // 2. Copy .dat files to base directory and archive the subdirs
    dateSubDirs.foreach(subDir => {
      val dirPath = subDir.getPath
      val datFiles = fs.listStatus(dirPath)
        .filter(status => !status.isDirectory && status.getPath.getName.endsWith(".dat"))
      
      // Copy each .dat file to base path (add date prefix to avoid filename conflicts)
      datFiles.foreach(file => {
        val sourcePath = file.getPath
        val targetFileName = s"${dirPath.getName}_${sourcePath.getName}"
        val targetPath = new Path(basePath, targetFileName)
        
        if (!fs.exists(targetPath)) {
          fs.copyFromLocalFile(false, true, sourcePath, targetPath)
          println(s"Successfully copied: ${sourcePath} -> ${targetPath}")
        } else {
          println(s"Skipping ${sourcePath} - target file ${targetFileName} already exists")
        }
      })

      // 3. Archive the subdirectory: Option 1 - Move to archive folder
      val archiveTarget = new Path(archiveRootPath, dirPath.getName)
      if (!fs.exists(archiveTarget)) {
        fs.rename(dirPath, archiveTarget)
        println(s"Archived directory: ${dirPath} -> ${archiveTarget}")
      } else {
        println(s"Skipping archive for ${dirPath} - target archive dir already exists")
      }

      // Alternative Option 2 - Create Hadoop Archive (HAR) for space efficiency
      // Uncomment this if you want compressed archives instead of moving directories
      // val harCommand = s"hadoop archive -archiveName ${dirPath.getName}.har -p ${basePath} ${dirPath.getName} ${archiveRootPath}"
      // val process = new ProcessBuilder(harCommand.split(" ")).start()
      // process.waitFor()
      // println(s"Created HAR archive for ${dirPath}")
    })
    
    // Cleanup Spark session
    spark.stop()
  }
}
Key Notes & Best Practices
  • Filename Conflict Prevention: The code appends the subdirectory's date to each .dat file (e.g., test1.dat becomes 20190101_test1.dat). Remove the date prefix if you don't need this, but be aware of overwriting risks.
  • Copy vs Move: The example uses copyFromLocalFile (a copy operation). For a more efficient move, replace it with fs.rename(sourcePath, targetPath)—this avoids duplicating data on HDFS.
  • Archive Options:
    • Move to Archive Folder: Simple, fast, and preserves the original directory structure. Great for basic cleanup.
    • HAR Archive: Use this if you want to reduce HDFS metadata overhead (useful for large numbers of small files). Ensure the hadoop command is available in your Spark executor's PATH if using this option.
  • Permissions: Confirm the Spark job has read/write permissions on /user/srav/ and the archive directory before running.
  • Cluster Deployment: When running on YARN, the HDFS config is automatically propagated to executors—no extra setup needed for standard clusters.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:27:41