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

SparkSQL读取多Parquet文件时添加子文件夹列的方法

当然可以!完全不用逐个加载文件就能搞定这个需求——Spark其实内置了专门的函数来帮你获取每条数据对应的源文件路径,再简单做个字符串处理就能提取到你要的子文件夹名称了。

核心思路

Spark提供了input_file_name()这个内置函数,它可以为DataFrame中的每一行返回其对应的源Parquet文件的完整路径。我们只需要基于这个路径,用字符串处理函数提取出目标子文件夹名称,再把它作为新列添加到DataFrame里就可以了,全程不需要逐个加载文件。

具体实现代码

假设你的Parquet文件路径格式是.../subfolderX/xxx.parquet,子文件夹是路径中倒数第二个部分,那么可以这样写:

// 保持你原来的批量加载方式不变
val originalDf = sqlContext.read.schema(schema).parquet(paths: _*)

// 添加subfolder列:拆分路径并提取倒数第二个部分(子文件夹名称)
val finalDf = originalDf.withColumn("subfolder",
    element_at(split(input_file_name(), "/"), -2)
)

灵活调整提取逻辑

如果你的路径结构不一样(比如子文件夹不是倒数第二个部分),可以用正则表达式来更精准地匹配:
比如你的路径是/data/project/subfolder1/my_file.parquet,想提取subfolder1,可以用regexp_extract:

val finalDf = originalDf.withColumn("subfolder",
    regexp_extract(input_file_name(), ".*/(subfolder\\d+)/.*", 1)
)

这里的正则表达式.*/(subfolder\\d+)/.*会匹配并捕获subfolder加数字的部分,第二个参数1表示取第一个捕获组的内容。

为什么这个方法更好

  • 不需要逐个加载文件,保持了批量加载的高效性
  • Spark会自动关联每条记录和它的源文件,不需要手动维护映射关系
  • 代码简洁易维护,调整提取逻辑只需要修改字符串处理部分就可以

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 12:37:42