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
相关产品推荐
相关产品推荐

