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

Spark Scala实现:获取文件路径及数据、给记录追加子目录名

嘿,这两个Spark Scala的问题都是日常数据处理里很常见的场景,我来一步步给你讲清楚怎么解决👇


问题1:在Spark Scala的RDD/DataFrame中获取文件路径及对应数据

用RDD API实现

有两种常用方案,适配不同的文件场景:

  • 小文件场景:用wholeTextFiles
    这个API直接返回「文件路径-文件完整内容」的键值对RDD,上手非常快:
import org.apache.spark.SparkContext
val sc: SparkContext = ... // 初始化你的SparkContext实例
val fileRDD = sc.wholeTextFiles("/path/to/target/directory")

// 打印每个文件的路径和内容
fileRDD.foreach { case (path, content) =>
  println(s"文件路径:$path\n内容:\n$content")
}

// 如果需要按行拆分内容,加个flatMap处理就行
val lineWithPathRDD = fileRDD.flatMap { case (path, content) =>
  content.split("\n").map(line => (path, line))
}
  • 大文件场景:textFile + 分区文件信息
    大文件用textFile更高效,我们可以通过mapPartitions获取每个分区对应的文件路径,再和每行数据关联:
import org.apache.spark.rdd.RDD
import org.apache.hadoop.mapreduce.lib.input.FileSplit
import org.apache.spark.TaskContext

val textRDD = sc.textFile("/path/to/target/directory/*")
val lineWithPathRDD: RDD[(String, String)] = textRDD.mapPartitions { iter =>
  // 获取当前分区对应的文件信息
  val split = TaskContext.get().getPartitionInfo.asInstanceOf[FileSplit]
  val filePath = split.getPath.toString
  // 把每行数据和文件路径绑定
  iter.map(line => (filePath, line))
}

用DataFrame/DataSet API实现

DataFrame有内置的input_file_name()函数,能直接获取每行数据对应的文件路径,代码更简洁:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.input_file_name

val spark: SparkSession = SparkSession.builder().appName("GetFilePathDemo").getOrCreate()
// 读取目录下的所有文本文件
val rawDF = spark.read.text("/path/to/target/directory")
// 添加文件路径列
val dfWithPath = rawDF.withColumn("file_path", input_file_name())

// 查看结果
dfWithPath.show(false)

结果里会有value列(每行原始内容)和file_path列(对应的文件路径)。


问题2:给每条记录追加所属子目录名称

这个需求可以基于上面的文件路径获取逻辑,先拿到路径再提取子目录名,最后拼接成目标格式。这里推荐用DataFrame实现,代码更易维护:

完整实现代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{input_file_name, regexp_extract, concat_ws}

val spark = SparkSession.builder()
  .appName("AppendSubdirNameDemo")
  .master("local[*]") // 本地测试用,生产环境请移除
  .getOrCreate()

// 读取所有子目录下的file1文件
val rawDF = spark.read.text("/parent_dir/*/file1")

// 核心逻辑:获取路径 → 提取子目录名 → 拼接结果
val resultDF = rawDF
  .withColumn("file_path", input_file_name())
  // 用正则从路径中提取子目录名,适配/parent_dir/subdirX/file1的格式
  .withColumn("subdir_name", regexp_extract($"file_path", ".*/parent_dir/([^/]+)/.*", 1))
  // 用|拼接原始内容和子目录名
  .withColumn("result", concat_ws("|", $"value", $"subdir_name"))
  .select("result")

// 打印结果验证
resultDF.show(false)

// 如果要保存结果到文件,执行下面代码(coalesce(1)可选,用来生成单个输出文件)
// resultDF.coalesce(1).write.text("/path/to/output/directory")

代码说明

  • regexp_extract的正则表达式:.*/parent_dir/([^/]+)/.* 精准匹配parent_dir后的子目录名,把它提取为subdir_name列
  • concat_ws("|", ...) 负责按要求的分隔符拼接内容,完美符合你要的输出格式

如果更习惯用RDD实现,逻辑类似:

import org.apache.spark.SparkContext
import org.apache.hadoop.mapreduce.lib.input.FileSplit
import org.apache.spark.TaskContext

val sc: SparkContext = ...
val textRDD = sc.textFile("/parent_dir/*/file1")
val resultRDD = textRDD.mapPartitions { iter =>
  val split = TaskContext.get().getPartitionInfo.asInstanceOf[FileSplit]
  val filePath = split.getPath.toString
  // 用字符串分割提取子目录名(也可以用正则)
  val pathParts = filePath.split("/")
  val subdirName = pathParts(pathParts.indexOf("parent_dir") + 1)
  // 拼接每行内容和子目录名
  iter.map(line => s"$line|$subdirName")
}

// 打印或保存结果
resultRDD.foreach(println)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:17:30