Databricks Autoloader抛IllegalArgumentException:示例为何缺schemaLocation?
尝试运行Databricks官网提供的最简Auto Loader示例,代码如下:
df = (spark.readStream.format("cloudFiles") .option("cloudFiles.format", "json") .load(input_data_path)) (df.writeStream.format("delta") .option("checkpointLocation", chkpt_path) .table("iot_stream"))
但持续收到报错:
IllegalArgumentException: cloudFiles.schemaLocation Could not find required option: schemaLocation. Please provide a schema location using
cloudFiles.schemaLocationfor storing inferred schema and supporting schema evolution.
疑问:如果cloudFiles.schemaLocation是必填项,为何各处示例均未包含该选项?底层原因是什么?
版本差异导致规则变化:Auto Loader在Databricks Runtime 9.0及以后的版本中,只要依赖自动Schema推断(未显式指定Schema),
cloudFiles.schemaLocation就会成为必填项。而官网的示例可能基于更早的Runtime版本(8.x及以前)编写——旧版本中,即使自动推断Schema也无需强制指定该路径,示例因此没有包含。示例隐含了显式指定Schema的前提:很多公开示例默认使用者会提前定义好数据Schema,这种场景下不需要
cloudFiles.schemaLocation。比如显式指定Schema的写法:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 提前定义数据结构 schema = StructType([ StructField("device_id", StringType()), StructField("temperature", IntegerType()), StructField("timestamp", StringType()) ]) df = (spark.readStream.format("cloudFiles") .option("cloudFiles.format", "json") .schema(schema) # 显式指定Schema,无需自动推断 .load(input_data_path))
此时Auto Loader直接使用给定的Schema,不需要持久化推断结果,自然不需要指定cloudFiles.schemaLocation。
- Auto Loader的Schema处理逻辑要求:当开启自动Schema推断时,
cloudFiles.schemaLocation是Auto Loader实现Schema演进、保证流处理一致性的核心路径。它会把首次推断的Schema存储在该路径,后续新文件流入时,会对比已有Schema并自动兼容新增字段(依据cloudFiles.schemaEvolutionMode配置)。如果缺少这个路径,Auto Loader无法持久化Schema信息,既无法处理Schema变化,也无法避免重复推断,因此会抛出必填项缺失的错误。
内容的提问来源于stack exchange,提问作者hikizume

