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

