Databricks Autoloader Schema Evolution抛出StateSchemaNotCompatible异常求助
核心问题
你遇到的StateSchemaNotCompatible异常,根源在于流处理中使用了distinct()算子。distinct()需要维护状态来追踪已处理过的记录,这个状态的Schema是基于初始流的Schema生成并存在checkpoint中的。当新字段触发Schema演化后,流的Schema发生了变化,但checkpoint里的状态Schema还是旧版本,两者不匹配,重试时就会抛出这个错误。
另外,你未设置cloudFiles.schemaEvolutionMode参数,默认情况下Autoloader遇到新字段会抛出NEW_FIELDS_IN_RECORD_WITH_FILE_PATH异常,这导致第一次任务失败,后续重试又触发了状态Schema不兼容的问题。
解决步骤
移除流中的
distinct()算子
流处理中的状态算子(如distinct、聚合、join等)的状态Schema无法随流Schema演化而更新,所以如果要开启Schema演化,就不能在流中使用这类算子。如果需要去重,建议改用以下方式:- 写入Delta表后,定期用批处理任务执行去重(比如
DELETE FROM table WHERE id IN (SELECT id FROM table GROUP BY id HAVING COUNT(*) > 1)); - 改用Delta的Merge操作实现幂等写入,基于唯一键判断是否插入新数据,替代流中的
distinct。
- 写入Delta表后,定期用批处理任务执行去重(比如
配置自动Schema演化参数
在读取流时添加cloudFiles.schemaEvolutionMode参数,设置为addNewColumns,这样遇到新字段时会自动更新Schema,而不是抛出异常:self.spark \ .readStream \ .format("cloudFiles") \ .option("cloudFiles.format", "json") \ .option("cloudFiles.inferColumnTypes", "true") \ .option("cloudFiles.schemaEvolutionMode", "addNewColumns") \ # 新增自动演化参数 .option("cloudFiles.schemaLocation", f"{self.target_s3_bucket}/_schema/{source_table_name}") \ .load(f"{self.source_s3_bucket}/{source_table_name}") \ # 移除distinct()算子 .writeStream \ .trigger(availableNow=True) \ .format("delta") \ .option("mergeSchema", "true") \ .option("checkpointLocation", f"{self.target_s3_bucket}/_checkpoint/{source_table_name}") \ .option("streamName", source_table_name) \ .start(f"{self.target_s3_bucket}/{target_table_name}")若必须保留
distinct()
这种情况下无法同时支持Schema演化,除非你能接受在Schema变化时删除旧的checkpoint目录,重新启动流(但这样会丢失之前的状态,可能导致重复数据)。
内容的提问来源于stack exchange,提问作者Robert Kossendey

