Spark加载多JSON文件并添加文件名列的高效实现方法
解决方案:为Spark DataFrame添加数据来源文件名列
不需要逐个循环文件,Spark提供了内置函数可以轻松实现这个需求,以下是具体实现步骤:
方法:使用input_file_name()内置函数
Spark SQL的input_file_name()函数可以直接获取每条数据对应的输入文件完整路径,之后可以按需提取纯文件名。
完整代码示例
from pyspark.sql.functions import input_file_name, regexp_extract # 读取文件夹下的所有多行JSON文件 data = spark.read.option("multiline", "true").json("/path/my_folder") # 添加包含完整文件路径的列 data_with_path = data.withColumn("full_file_path", input_file_name()) # 提取纯文件名(如aaa.json、bbb.json),这里用正则表达式截取路径最后一段 data_with_filename = data_with_path.withColumn( "filename", regexp_extract("full_file_path", ".*/(.*)", 1) ) # 也可以用split+element_at的方式提取文件名(更直观) # from pyspark.sql.functions import split, element_at # data_with_filename = data_with_path.withColumn( # "filename", # element_at(split("full_file_path", "/"), -1) # ) # 查看结果 data_with_filename.select("first_name", "last_name", "filename").show()
说明
input_file_name():返回当前数据行对应的源文件完整路径,比如/path/my_folder/aaa.json- 正则提取或
split+element_at:用于从完整路径中剥离出单纯的文件名,两种方法都能达到目的,可根据个人习惯选择
内容的提问来源于stack exchange,提问作者Liky
相关产品推荐
相关产品推荐

