开启CDF的Delta Lake执行Vacuum后读取变更数据报文件未找到异常
Delta Lake CDC读取遇文件找不到错误的原因及解决办法
问题根源
你通过以下代码写入带变更数据捕获(CDC)的Delta表:
df.write.format("delta").partitionBy("g","p").option("delta.enableChangeDataFeed", "true").mode("append").save(path)
在版本3、4插入数据,版本5删除部分数据后,执行了deltaTable.vacuum(8)——这个命令会物理删除8小时前生成的所有旧数据文件。
当你尝试从版本3开始读取CDC数据时:
spark.read.format("delta") .option("readChangeFeed", "true") .option("startingVersion", 3) .load(path)
Delta需要先加载版本3对应的基础parquet文件来构建初始数据快照,再结合版本4、5的变更日志生成完整的变更流。如果版本3的文件生成时间距离vacuum执行已经超过8小时,这些文件已经被GCS彻底删除,自然会触发FileNotFoundException。删除集群重试也没用,因为文件是存储在GCS上的,和集群无关。
解决办法
- 调整vacuum保留时长:如果需要回溯读取旧版本的CDC数据,
vacuum的retentionHours参数必须大于「起始版本生成时间到当前时间」的间隔,确保起始版本的文件不会被清理。比如要保留7天的文件,执行:deltaTable.vacuum(168) - 恢复已删除文件:如果文件已经被vacuum清理,只能通过GCS的版本控制功能(如果已开启)恢复被删除的文件,或者从更早的备份重新生成数据。
- CDC与vacuum的配合规则:启用CDC后,必须保证vacuum的保留策略不会清理掉你需要读取的起始版本对应的基础文件,否则CDC读取必然失败。
内容的提问来源于stack exchange,提问作者Bindu
相关产品推荐
相关产品推荐

