PySpark指定Schema读JSON时使用input_file_name()返回空值问题
问题根因
你使用的Spark 2.3.2版本存在JSON数据源的已知Bug:当使用自定义Schema读取JSON文件,且查询中包含input_file_name()这类文件级伪列时,会触发字段解析偏移错误,导致业务字段返回null。
具体触发逻辑
- 自定义Schema读取时,Spark默认开启JSON列剪枝优化,只会读取Schema中声明的字段来提升性能。当查询包含
input_file_name()伪列时,Spark会错误地将该伪列插入到待解析字段列表的最前端,导致后续JSON业务字段的解析偏移量全部错位,字段值被判定为不匹配Schema,最终返回null。 - 仅查询业务字段不包含伪列时,没有额外字段插入待解析列表,偏移量匹配,所以可以正常解析出id的值。
- 不指定Schema使用自动推断时,Spark会先扫描全量数据生成完整Schema,列剪枝逻辑不会将伪列插入原始字段列表头部,偏移量正常,因此不会出现解析失败问题。
解决方案
临时规避方案(无需改版本/配置)
读取数据后先缓存避免重复读取,再通过withColumn新增文件名列,不要在首次select时直接调用input_file_name():
# 按自定义Schema读取数据 mydf=spark.read.schema(myschema).option("multiline","true").option("mode","PERMISSIVE").json("/user/testuser/test1.json") # 缓存后新增文件名列 mydf = mydf.cache() mydf = mydf.withColumn("file_path", input_file_name()) # 正常查询即可 mydf.select("id", "file_path").show()
配置修改方案
关闭JSON列剪枝优化即可解决偏移错误问题,性能会有小幅下降:
spark.conf.set("spark.sql.json.enableColumnPruning", "false")
永久修复方案
将Spark版本升级到2.4.0及以上,该Bug已经在Spark 2.4.0版本被官方修复。
内容的提问来源于stack exchange,提问作者Seaport
相关产品推荐
相关产品推荐

