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

Databricks流处理DeltaFileNotFoundException问题求助(需保留Checkpoint)

问题

使用Databricks进行流处理开发,现有两个作业(notebook1和notebook2)从不同数据源读取数据,以标准格式写入同一本地数据集(LDS),通过分区隔离避免冲突,该方案已稳定运行近5个月。今日出现错误:

com.databricks.sql.transaction.tahoe.DeltaFileNotFoundException: No file found in the directory: dbfs:/mnt/streaming/streaming1/_delta_log.

已找到两种方案,但存在疑问:

  • 方案1:建议使用新的checkpoint目录,或在集群Spark配置中设置spark.sql.files.ignoreMissingFiles为true。但新checkpoint会导致全量重跑,不符合需求,想了解设置该属性后,是从上次checkpoint恢复处理,还是从头开始?
  • 方案2:提到修改父目录并在start()中添加目录选项,完全无法理解,且已有类似配置,似乎无法解决问题。

流处理核心代码:

spark.readStream.format("delta") \
  .option("readChangeFeed", "true") \
  .option("maxFilesPerTrigger", 250) \
  .option("maxBytesPerTrigger", 536870912)\
  .option("failOnDataLoss", "true")\
  .load(DATA_PATH)\
  .filter(expr("_change_type not in ('delete', 'update_preimage')"))\
  .writeStream\
  .queryName(streamQueryName)\
  .foreachBatch(MainFunctionstoprocess)\
  .option("checkpointLocation", checkpointLocation)\
  .option("mergeSchema", "true")\
  .trigger(processingTime='1 seconds')\
  .start()

需求:如何在不删除checkpoint的情况下解决该问题,恢复到失败前的处理位置,或回到某个checkpoint仅重跑部分数据?


解决方案

关于spark.sql.files.ignoreMissingFiles的恢复逻辑

设置spark.sql.files.ignoreMissingFiles为true后,流作业会从上次checkpoint记录的位置恢复处理,不会从头开始。这个配置的作用是让Spark忽略读取过程中遇到的缺失文件,不会因为单个日志文件丢失就中断作业,而是跳过缺失部分,继续处理后续存在的日志和数据。

具体解决步骤

  1. 配置Spark属性
    可以在Databricks集群的Spark配置中添加全局配置:

    spark.sql.files.ignoreMissingFiles true
    

    也可以在作业代码开头动态设置:

    spark.conf.set("spark.sql.files.ignoreMissingFiles", "true")
    
  2. 调整failOnDataLoss参数
    代码中当前设置failOnDataLoss为true,这会导致作业检测到数据丢失时直接失败。结合忽略缺失文件的配置,建议将该参数改为false,避免因日志文件缺失触发失败:

    spark.readStream.format("delta") \
      .option("readChangeFeed", "true") \
      .option("maxFilesPerTrigger", 250) \
      .option("maxBytesPerTrigger", 536870912)\
      .option("failOnDataLoss", "false")\  # 修改此处
      .load(DATA_PATH)\
      # 后续代码保持不变
    
  3. 重启流作业
    无需删除现有checkpoint,直接重启作业即可。作业会读取checkpoint中记录的偏移量,跳过缺失的日志文件,从上次成功处理的位置继续运行。

重跑部分数据的操作

如果需要回到特定checkpoint位置重跑:

  1. 找到目标checkpoint对应的偏移量记录(在checkpoint目录下的offsets子目录中,以JSON文件存储)。
  2. 备份当前checkpoint目录,将目标偏移量文件替换到offsets目录中。
  3. 按照上述步骤配置好参数后,重启作业,作业会从替换后的偏移量位置开始处理。

注意事项

  • 先确认dbfs:/mnt/streaming/streaming1/_delta_log目录是否存在部分文件,若目录完全为空,需检查数据源Delta表是否被误删除,这种情况下即使配置忽略缺失文件,作业也可能无法正常运行。
  • 确保两个作业的分区逻辑严格隔离,避免因分区重叠导致日志文件冲突或丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 06:24:35