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

Spark/Scala实现按国家处理文件夹文件并迁移到对应层级目录

Spark Scala 按国家维度处理文件并分目录存储实现方案

实现思路

  • 扫描源路径下所有符合命名规则的csv文件
  • 从文件名中提取国家编码,年月字段可选择从文件名提取或者自定义目标年月
  • 按国家编码对文件分组,对每个国家的全量文件执行自定义处理逻辑
  • 处理完成后将结果写入年/月/国家编码层级的目标路径,可按需选择覆盖、追加等写入模式

文件名提取规则适配示例格式Casedata_${国家码}_${年月}_${时间戳}.csv,用下划线分割后第2位为3位国家编码,第3位为6位年月字段,命名规则有变动可自行调整提取逻辑


完整代码实现

Spark Scala版本(适合大文件量、分布式集群场景)

import org.apache.spark.sql.SparkSession
import java.io.File

object CountryFileProcess {
  def main(args: Array[String]): Unit = {
    // 初始化SparkSession,集群运行时去掉master配置
    val spark = SparkSession.builder()
      .appName("CountryFileProcess")
      .master("local[*]")
      .getOrCreate()
      
    // 可修改配置项
    val sourceDir = "/your/source/file/dir" // 源文件根目录
    val targetRootDir = "/your/target/root/dir" // 目标路径根目录
    val useSourceFileMonth = false // 设为true则从源文件名取年月,false则用下方自定义目标年月
    val customTargetMonth = "202111" // 示例目标年月
    
    // 自定义文件处理逻辑,此处示例为读取csv去重,可替换为实际业务逻辑
    def processCsvFile(filePath: String) = {
      spark.read
        .option("header", "true")
        .csv(filePath)
        .dropDuplicates()
    }
    
    // 扫描源目录下符合命名规则的csv文件
    val validSourceFiles = new File(sourceDir).listFiles()
      .filter(_.getName.matches("Casedata_[A-Z]{3}_\\d{6}_.*\\.csv"))
      
    // 按国家编码分组
    val countryGroupFiles = validSourceFiles.groupBy { file =>
      file.getName.split("_")(1)
    }
    
    // 逐个国家处理写入
    countryGroupFiles.foreach { case (countryCode, files) =>
      // 构造目标路径
      val finalMonth = if (useSourceFileMonth) files.head.getName.split("_")(2) else customTargetMonth
      val year = finalMonth.substring(0,4)
      val month = finalMonth.substring(4,6)
      val targetPath = s"$targetRootDir/$year/$month/$countryCode"
      
      // 合并当前国家所有文件并处理
      val countryTotalDf = files.map(f => processCsvFile(f.getAbsolutePath)).reduce(_ union _)
      
      // 写入目标路径,coalesce(1)表示输出单个文件,不需要可删除
      countryTotalDf.coalesce(1)
        .write
        .option("header", "true")
        .mode("overwrite") // 写入模式可选append/ignore/errorifexists
        .csv(targetPath)
    }
    spark.stop()
  }
}

原生Scala版本(适合小文件量、本地运行场景)

import java.io.{File, FileInputStream, FileOutputStream}
import java.nio.channels.FileChannel

object LocalCountryFileProcess {
  def main(args: Array[String]): Unit = {
    val sourceDir = "/your/source/file/dir"
    val targetRootDir = "/your/target/root/dir"
    val useSourceFileMonth = false
    val customTargetMonth = "202111"
    
    // 自定义文件处理逻辑,此处示例为直接拷贝文件,可替换为读取内容加工后写入
    def transferFile(source: File, target: File): Unit = {
      val inputChannel = new FileInputStream(source).getChannel
      val outputChannel = new FileOutputStream(target).getChannel
      outputChannel.transferFrom(inputChannel, 0, inputChannel.size())
      inputChannel.close()
      outputChannel.close()
    }
    
    val validFiles = new File(sourceDir).listFiles()
      .filter(_.getName.matches("Casedata_[A-Z]{3}_\\d{6}_.*\\.csv"))
      
    val countryGroup = validFiles.groupBy(_.getName.split("_")(1))
    
    countryGroup.foreach { case (countryCode, files) =>
      val finalMonth = if (useSourceFileMonth) files.head.getName.split("_")(2) else customTargetMonth
      val year = finalMonth.substring(0,4)
      val month = finalMonth.substring(4,6)
      val targetDir = new File(s"$targetRootDir/$year/$month/$countryCode")
      if (!targetDir.exists()) targetDir.mkdirs()
      
      files.foreach(file => {
        val targetFile = new File(targetDir, file.getName)
        transferFile(file, targetFile)
      })
    }
  }
}

注意事项

  • Spark写入后生成的文件默认是part开头的命名,需要自定义文件名可在写入完成后用原生Scala文件操作重命名
  • 文件名正则匹配规则可根据实际命名格式调整,当前规则适配3位大写字母国家码、6位数字年月的格式
  • 处理大文件时建议使用Spark版本,避免本地内存溢出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 07:24:05