写入Delta表时出现校验和错误,求解决方法(Spark3.4/Delta2.4)
问题背景
环境与表信息
- Spark版本:3.4(开源版)
- Delta Lake版本:2.4
- 表总记录数:15602
- Delta表总版本数:15602
- Delta表大小:1.4 GB
- 分区文件夹数量:488
报错日志(翻译后)
24/01/27 13:42:49 ERROR Executor: 阶段688.0中的任务1.0(TID 819)执行异常 org.apache.spark.SparkException:
读取文件file:///home/kotesh/delete_data_compac/my_table/_delta_log/00000000000000015603.json时遇到错误。
详情:org.apache.spark.sql.errors.QueryExecutionErrors$.cannotReadFilesError(QueryExecutionErrors.scala:877)
...(中间调用栈省略)
引起原因:org.apache.hadoop.fs.ChecksumException: 校验和错误:file:/home/kotesh/delete_data_compac/my_table/_delta_log/00000000000000015603.json
位置0处预期校验和:-70964045,实际得到:69110470
at org.apache.hadoop.fs.FSInputChecker.verifySums (FSInputChecker.java:347)
该错误属于Hadoop文件系统层面的校验失败,指向Delta日志目录下的00000000000000015603.json文件损坏或校验文件异常。
解决方法
1. 先确认文件状态
手动检查损坏的日志文件是否可读:
cat /home/kotesh/delete_data_compac/my_table/_delta_log/00000000000000015603.json
- 如果文件无法打开、内容乱码,说明文件本身损坏;
- 若内容正常,仅同目录下的
.00000000000000015603.json.crc校验文件异常,可直接删除该校验文件后重试任务。
2. 临时应急:删除损坏的最新日志
由于表当前版本为15602,损坏的是未提交的15603版本日志,直接删除该日志及对应校验文件:
rm /home/kotesh/delete_data_compac/my_table/_delta_log/00000000000000015603.json rm /home/kotesh/delete_data_compac/my_table/_delta_log/.00000000000000015603.json.crc
之后重新执行插入操作,Delta会自动生成新的日志文件。
3. 回滚到正常版本
如果损坏的日志对应已提交的版本,执行回滚命令恢复到上一个正常版本(15602):
Spark SQL方式:
RESTORE TABLE my_table TO VERSION AS OF 15602;
Scala API方式:
import io.delta.tables._ val deltaTable = DeltaTable.forPath(spark, "/home/kotesh/delete_data_compac/my_table") deltaTable.restoreToVersion(15602)
回滚完成后重新执行插入操作。
4. 重建表(多日志损坏时)
若多个日志文件损坏,建议导出有效数据后重建表:
- 导出正常版本的数据到临时表:
CREATE TABLE temp_table AS SELECT * FROM my_table VERSION AS OF 15602;
- 删除原Delta表目录:
rm -rf /home/kotesh/delete_data_compac/my_table
- 重新创建Delta表并导入数据:
CREATE TABLE my_table USING DELTA AS SELECT * FROM temp_table;
- 清理临时表:
DROP TABLE temp_table;
预防措施
- 定期执行合并操作减少日志文件数量:
OPTIMIZE my_table; - 禁止手动修改Delta表
_delta_log目录下的任何文件 - 确保存储介质(本地磁盘/分布式存储)稳定,避免文件损坏
内容的提问来源于stack exchange,提问作者kotesh

