You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 14:54:09