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

修改集群配置后Databricks流任务报文件缺失,求配置添加位置

问题描述

此前Spark Structured Streaming脚本运行正常,修改Databricks集群配置并通过Keyvault scope挂载存储位置后,执行脚本时出现文件缺失错误:

"Error while reading file
/mnt/topics/audit-logs/minio.sys.tmp/multipart/v1/c782.025a337657/azure.json.
A file notification was received for file:
/mnt/topics/audit-logs/minio.sys.tmp/multipart/v1/c782.025a337657/azure.json
but it does not exist anymore. Please ensure that files are not
deleted before they are processed. To continue your stream, you can
set the Spark SQL configuration spark.sql.files.ignoreMissingFiles to
true."

执行的脚本如下:

table_name= "main.deltalake_db.Table_test"
checkpoint_path = "/mnt/_checkpoint"
file_path ="/mnt/topics/audit-logs"
schema = "authenticationMechanism STRING,authenticationMechanismId STRING"
(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .schema(schema)  
  .load(file_path)
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable(table_name))

即使改回原集群配置,问题仍然存在。求助:spark.sql.files.ignoreMissingFiles配置应添加到脚本的哪个位置?

解决方案

spark.sql.files.ignoreMissingFiles是Spark SQL的全局配置项,有两种添加方式:

方式1:脚本开头全局设置

在脚本最开始通过spark.conf.set()配置,作用于整个SparkSession:

# 全局设置忽略缺失文件
spark.conf.set("spark.sql.files.ignoreMissingFiles", "true")

table_name= "main.deltalake_db.Table_test"
checkpoint_path = "/mnt/_checkpoint"
file_path ="/mnt/topics/audit-logs"
schema = "authenticationMechanism STRING,authenticationMechanismId STRING"
(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .schema(schema)  
  .load(file_path)
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable(table_name))

方式2:在readStream阶段单独设置

在readStream链式调用中通过.option()添加,仅作用于当前流读取任务:

table_name= "main.deltalake_db.Table_test"
checkpoint_path = "/mnt/_checkpoint"
file_path ="/mnt/topics/audit-logs"
schema = "authenticationMechanism STRING,authenticationMechanismId STRING"
(spark.readStream
  .format("cloudFiles")
  .option("cloudFiles.format", "json")
  .option("cloudFiles.schemaLocation", checkpoint_path)
  .option("spark.sql.files.ignoreMissingFiles", "true")  # 单独为当前流配置
  .schema(schema)  
  .load(file_path)
  .writeStream
  .option("checkpointLocation", checkpoint_path)
  .trigger(availableNow=True)
  .toTable(table_name))

注意事项

  • 推荐全局配置方式,避免后续其他流任务遇到同类问题
  • 配置值需设为字符串"true",而非布尔值True,否则可能不生效
  • 改回原集群配置仍报错,大概率是之前流任务的checkpoint记录了已删除文件的元数据,添加该配置后可跳过这些无效记录,让流任务继续运行

内容的提问来源于stack exchange,提问作者ZZZSharePoint

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 10:35:25