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

