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

Databricks Autoloader使用cloudFiles时遇FileNotFoundException问题排查

Spark Structured Streaming 增量摄入报错分析

代码示例

val ssdf = spark.readStream.schema(schema) 
.format("cloudFiles")
.option("cloudFiles.format", "parquet")
.load(sourceUrl)
.filter(criteriaFilter)

val transformedDf = ssdf.transform(.....)

val processData = transformedDf
.select(recordFields: _*) 
.writeStream
.option("checkpointLocation", outputUrl + "checkpoint/")
.format("parquet")
.outputMode("append")
.option("path", outputUrl + run_id + "/") 
.trigger(Trigger.Once())
.start()
processData.processAllAvailable()
processData.stop()

问题场景

该代码用于每周增量摄入数据,每周数据写入带run_id的新文件夹,检查点固定存储在outputUrl/checkpoint/。前3-5次运行正常,近期触发如下错误:

ERROR: Query termination received for [id=2345245425], with exception: org.apache.spark.SparkException: Job aborted.
Caused by: java.io.FileNotFoundException: Unable to find batch s3://outputPath/20230810063959/_spark_metadata/0

问题原因

  • 检查点与输出目录不匹配:使用Trigger.Once()时,Spark会在输出目录生成_spark_metadata文件夹,记录批次ID、偏移量等元数据。但你的任务每次输出到新目录,检查点却固定不变,导致检查点中存储的批次信息会指向旧的输出目录。如果旧目录被清理(比如删除历史数据),后续任务读取检查点时,就会找不到对应目录下的_spark_metadata文件。
  • 批次ID状态残留:Spark的批次ID从0开始递增,检查点会持久化历史批次状态。当旧输出目录被删除后,新任务启动时会尝试根据检查点中的批次ID去对应目录查找元数据文件,而该目录已不存在,最终触发文件未找到的错误。

为什么会在当前文件夹查找0文件

当新任务启动时,Spark会优先从检查点读取上一次的运行状态。如果检查点中记录的批次ID对应的旧输出目录已被删除,或者检查点元数据出现异常,Spark会尝试在当前任务的新输出目录中查找历史批次的元数据文件。但新目录是刚创建的,没有之前批次(比如ID为0的批次)的_spark_metadata/0文件,因此抛出该错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 10:42:15