Databricks Autoloader启用rescue模式加载Parquet schema报错排查
我有一批Schema持续演进的Parquet文件,需要全部加载到单个Delta表中。计划使用Autoloader并设置schemaEvolutionMode="rescue"(让所有与目标Schema不匹配的源字段落入_rescued_data列),同时给Autoloader指定了.schema(target_schema)。但读取部分文件时出现报错:
Invalid Spark read type: expected optional group my_column (LIST) { repeated group list { optional binary element (STRING); } } to be list but found Some(StringType)
目标表中my_column的数据类型为String,为什么该字段没有被导入_rescued_data列反而触发报错?
使用的代码如下:
read_options = { "cloudFiles.format": "parquet", "cloudFiles.schemaLocation": "some location", "cloudFiles.schemaEvolutionMode": "rescue" } spark.readStream.format("cloudFiles") .options(**read_options) .schema(target_schema) .load("source_path") .foreachBatch(<save function>) .outputMode("append") .trigger("availableNow", True) .start()
使用的Databricks版本为13.2(Spark 3.4.0,Scala 2.12)
核心原因
schemaEvolutionMode="rescue"仅处理新增字段或可兼容类型转换的场景,当源数据中已有字段的类型与指定的target_schema类型完全不兼容且无默认转换逻辑时(比如这里的LIST数组类型和String类型),Autoloader不会将该字段放入_rescued_data,而是直接抛出类型不匹配错误。Spark的类型规则里,数组和String属于完全不兼容类型,没有自动转换路径,因此触发报错。
解决方案
调整目标Schema并添加类型提示
如果业务允许,将目标表my_column的类型改为Array[String],同时在读取参数中添加schemaHints明确字段类型,让源数据的LIST类型正常匹配,后续String类型的该字段也能通过提示规则处理:read_options = { "cloudFiles.format": "parquet", "cloudFiles.schemaLocation": "some location", "cloudFiles.schemaEvolutionMode": "rescue", "cloudFiles.schemaHints": "my_column ARRAY<STRING>" }自定义批次处理逻辑
移除.schema(target_schema)的指定,让Autoloader自动推断源Schema,然后在foreachBatch的保存函数中手动处理类型差异:- 对
my_column进行类型判断,若为数组则用concat_ws等函数转为String,若为String则直接保留; - 将无法处理的异常字段手动写入
_rescued_data列。
- 对
预处理源文件
提前对存在类型不匹配的Parquet文件进行转换,统一my_column的类型后再通过Autoloader加载。
内容的提问来源于stack exchange,提问作者archjkeee

