Databricks Autoloader/writeStream如何配置自动重试?
Databricks AutoLoader 处理UNKNOWN_FIELD_EXCEPTION.NEW_FIELDS_IN_FILE 错误方案
错误核心解析
这个错误提示里的「自动重试修复」是AutoLoader内置的schema演化机制,但需要配置正确才能触发,不需要额外手动配置重试逻辑,问题出在现有配置的冲突或缺失上。
配置问题排查与修正
CSV格式的schema推断开关
CSV属于无强schema格式,必须开启AutoLoader的schema推断才能识别新增列。你的代码里缺少cloudFiles.inferSchema配置,这会导致AutoLoader无法自动检测新列,进而触发该错误。Delta写入的schema配置冲突
同时设置mergeSchema和overwriteSchema会导致冲突——前者是合并新增列,后者是直接覆盖现有schema,建议只保留mergeSchema即可。schemaLocation的残留问题
如果dbfs:/mnt/temp/checkpoints/schema路径下已经存储了旧版本的schema文件,AutoLoader会优先使用旧schema,遇到新列时就会报错。可以先删除该路径下的所有文件,让AutoLoader重新扫描并生成包含新列的完整schema。
修正后的代码
xd = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "csv") \ .option("cloudFiles.schemaLocation", "dbfs:/mnt/temp/checkpoints/schema") \ .option("cloudFiles.schemaEvolutionMode", "addNewColumns") \ .option("cloudFiles.inferSchema", "true") \ .option("pathGlobfilter", "20230808_*") \ .load("/mnt/test_loc") \ .writeStream \ .format("delta") \ .outputMode("append") \ .option("checkpointLocation", "dbfs:/mnt/temp/checkpoints") \ .option("mergeSchema", "true") \ .toTable("dBronze.Table1")
额外注意事项
- 确保所有CSV文件的表头格式一致,新增列的位置不影响AutoLoader的识别,但必须保证表头存在。
- 检查
schemaLocation和checkpointLocation路径的读写权限,AutoLoader需要写入更新后的schema和 checkpoint 信息。 - 错误提示里的
automatic retry: true表示AutoLoader会自动重试读取文件并更新schema,但前提是配置允许schema演化。如果重试后仍报错,建议检查文件是否存在格式错误(比如分隔符不统一、数据类型不匹配)。
内容的提问来源于stack exchange,提问作者John Stud
相关产品推荐
相关产品推荐

