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
相关产品推荐
相关产品推荐

