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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:09:29