PySpark调用input_file_name()读取json.gz文件返回空字符串问题
PySpark调用input_file_name()获取json.gz文件名返回空字符串解决方案
问题复现
读取json.gz格式文件时,两种常规写法获取源文件名均返回空字符串:
- 提取路径末尾文件名的写法:
df.withColumn("source_file",sql_f.element_at(sql_f.split(sql_f.input_file_name(), "/"), -1))
- 直接调用函数的写法:
df.withColumn("source_file",sql_f.input_file_name())
问题原因
- 最核心原因:
input_file_name()是文件块级别的元数据函数,仅在未经过shuffle操作的初始读取DataFrame上生效。如果在调用该函数前,DataFrame已经执行过groupBy、join、repartition、distinct、聚合计算等会触发shuffle的算子,原始输入文件的元信息会被丢弃,函数必然返回空字符串。 - 次要原因:Spark 3.1以下版本对压缩文件的元数据传递存在兼容性问题,部分自定义读取逻辑也会导致元数据丢失。
可行解决方案
方案1:读取文件后立刻添加文件名列(适配所有Spark版本)
保证添加source_file列的操作是文件读取后的第一个转换操作,中间不要插入任何shuffle类算子,即可正常获取文件名。
正确示例:
from pyspark.sql import functions as sql_f # 第一步:读取json.gz文件,读取后不要做任何聚合、关联、重分区操作 df = spark.read.json("/your/json/gz/file/path/") # 第二步:立刻添加文件名列,此时元数据未丢失 df = df.withColumn("source_file", sql_f.element_at(sql_f.split(sql_f.input_file_name(), "/"), -1)) # 后续再执行清洗、聚合、关联等任意操作,都不会影响source_file列的取值
方案2:使用内置_metadata列(适配Spark 3.1+版本,稳定性更高)
Spark 3.1及以上版本提供持久化的_metadata内置列,文件元数据不会因为常规shuffle操作丢失,比input_file_name()容错性更强:
from pyspark.sql import functions as sql_f df = spark.read.json("/your/json/gz/file/path/") \ .withColumn("source_file", sql_f.element_at(sql_f.split(sql_f.col("_metadata.file_path"), "/"), -1))
可通过df.select("_metadata").show(truncate=False)查看所有可用文件元信息,包含文件路径、大小、修改时间等字段。
方案3:低版本Spark兼容方案
如果使用Spark 3.1以下版本,且需要在后续shuffle操作中保留文件名,可在添加文件名列后先缓存DataFrame,固定列值再做后续计算:
from pyspark.sql import functions as sql_f from pyspark.storagelevel import StorageLevel df = spark.read.json("/your/json/gz/file/path/") # 先添加文件名列再缓存 df = df.withColumn("source_file", sql_f.element_at(sql_f.split(sql_f.input_file_name(), "/"), -1)) df.persist(StorageLevel.MEMORY_AND_DISK) # 触发缓存执行,固定列值 df.count() # 后续执行任意shuffle操作,source_file列都不会返回空值
内容的提问来源于stack exchange,提问作者Nabeel Khan Ghauri
相关产品推荐
相关产品推荐

