Spark加载/合并多文件时给DataFrame加文件名列的方法及代码排查
关于Spark合并文件时添加文件名列的问题
嘿,我来帮你搞定这个问题!先拆解你代码里的问题,再给你正确的实现方式。
你的代码存在的两个核心问题
- Spark DataFrame是不可变对象:你调用
df.withColumn("file_name", input_file_name)后,这个操作会生成一个新的DataFrame,但你没有把这个新对象赋值给任何变量。原来的df还是保持原样,不会新增file_name列。 - 内置函数调用错误:
input_file_name是Spark的内置函数,需要加括号()来调用。直接写input_file_name会把函数本身当作值传入,而不是执行函数获取当前记录对应的文件名。
正确的实现代码
把这两个问题修正后,代码应该是这样的:
import org.apache.spark.sql.functions._ // 先加载原始数据 val rawDf = spark.read.format("csv") .option("sep","|") .option("inferSchema","true") .option("header","false") .load("mnt/rawdata/2019/01/01/corp/ABC*.gz") // 添加文件名列,并将结果赋值给新变量(或覆盖原变量) val dfWithFileName = rawDf.withColumn("file_name", input_file_name())
进阶:提取文件名的特定部分
如果你的需求是只保留文件名主体(比如去掉完整路径和.gz后缀),可以用regexp_extract函数处理:
val dfWithCleanFileName = rawDf.withColumn("file_name", // 正则匹配最后一个/之后的部分,提取.gz之前的文件名 regexp_extract(input_file_name(), """([^/]+)\.gz$""", 1) )
这样每条记录就会准确对应它的来源文件名啦~
内容的提问来源于stack exchange,提问作者ASH
相关产品推荐
相关产品推荐

