使用Autoloader无法访问部分JSON属性的排查求助
问题描述
同一个JSON文件被两个不同的Autoloader加载:
- 第一个启用Schema演进,仅替换JSON属性名中的空格后写入Delta表,所有值正常显示;
- 第二个映射至自定义Schema,仅使用部分属性,通过
withColumn处理后筛选列,但部分字段始终为Null(如status.name),其他字段(如risk_severity.value)正常取值。
Autoloader配置
df = (spark .readStream .format('cloudFiles') .option('cloudFiles.format', 'json') .option('multiLine', 'true') .option('cloudFiles.schemaEvolutionMode','rescue') .option('cloudFiles.includeExistingFiles','true') .option('cloudFiles.schemaLocation', bronze_schema) .option('cloudFiles.inferColumnTypes', 'true') .option('pathGlobFilter','*.json') .load(upload_path) .transform(lambda df: remove_spaces_from_columns(df)) .withColumn(...)
写入器配置
df.writeStream.format('delta') \ .queryName(al_stream_name) \ .outputMode('append') \ .option('checkpointLocation', checkpoint_path) \ .option('mergeSchema', 'true') \ .trigger(once = True) \ .table(bronze_table)
异常示例
.withColumn('vl_rating', col('risk_severity.value')) # 正常获取值 .withColumn('status', col('status.name')) # 始终为Null ... .select( 'rating', 'status', ...
离线加载该JSON文件时可正常获取
status.name的值,代码如下:
%python import pyspark.sql.functions as psf jsontest = spark.read.option('inferSchema','true').json('dbfs:....json') df = jsontest.withColumn('status', psf.col('status.name')).select('status') display(df)
排查思路
- 检查自定义Schema的匹配性:如果第二个Autoloader使用了自定义Schema,确认
status字段是否被定义为结构体类型。若初始Schema错误,即使开启Schema演进,也可能无法正确解析嵌套字段。 - 验证
remove_spaces_from_columns函数影响:临时注释该transform步骤,测试流处理是否能正常获取status.name。排查函数是否误修改了嵌套字段的名称结构,或对嵌套字段处理逻辑存在漏洞。 - 查看
_rescued_data列内容:由于开启了schemaEvolutionMode='rescue',流DataFrame会生成_rescued_data列,检查其中是否包含未被正确解析的status相关数据,判断是否是Schema推断/自定义Schema未覆盖正确结构。 - 清理Schema和Checkpoint缓存:Autoloader会缓存旧Schema到
schemaLocation路径,Checkpoint也会记录流状态。删除这两个路径下的文件,重新运行流任务,排除旧缓存导致的解析异常。 - 确认多线JSON配置一致性:确保流处理和离线加载都开启了
multiLine=true,避免流处理时因JSON换行问题导致部分字段解析失败。 - 检查字段大小写敏感性:查看
spark.sql.caseSensitive配置,若JSON中字段是Status.Name而非status.name,流处理中用小写引用会返回Null。打印流DataFrame的Schema,对比离线加载的Schema结构。 - 测试简化流任务:去掉所有
withColumn和select逻辑,直接将原始流DataFrame写入临时表,查看status字段的结构和值是否正常,逐步添加逻辑定位问题点。 - 排查数据源特殊性:检查流处理的数据源中是否存在部分JSON文件的
status结构异常(比如status是字符串而非对象),离线测试可能仅用了单个正常文件,而流处理包含异常文件。
内容的提问来源于stack exchange,提问作者Chris de Groot
相关产品推荐
相关产品推荐

