删除S3上Hudi数据集分区后出现FileNotFoundException的解决方法
问题描述
误删除S3上Hudi分区pt=2024-01后,尝试用相同分区名更新Hudi表时触发FileNotFoundException。已尝试用Spark清除该分区痕迹、删除.hoodie文件夹中该分区的提交文件等操作,但均未生效。
执行的Spark代码
val basePath = <<TablePath>> val df = spark.read.format("org.apache.hudi"). option("hoodie.datasource.read.extract.partition.values.from.path", true). load(basePath) val delDf = df.filter(<<Filter Condition>>) delDf. write. format("org.apache.hudi"). option(OPERATION_OPT_KEY,"delete"). option(PRECOMBINE_FIELD_OPT_KEY, <<Precombine Key>>). option(RECORDKEY_FIELD_OPT_KEY, <<Record Key>>). // 必要时添加分区键 option("hoodie.metrics.on", "false"). option("hoodie.write.tagged.record.storage.level", "DISK_ONLY"). option("hoodie.write.status.storage.level", "DISK_ONLY"). mode(Append). save(basePath)
报错日志(中文翻译)
用户类抛出异常: org.apache.hudi.exception.HoodieUpsertException: 提交时间20240112113537886的upsert操作失败 at org.apache.hudi.table.action.commit.BaseWriteHelper.write(BaseWriteHelper.java:64) at org.apache.hudi.table.action.commit.SparkUpsertCommitActionExecutor.execute(SparkUpsertCommitActionExecutor.java:45) at org.apache.hudi.table.HoodieSparkCopyOnWriteTable.upsert(HoodieSparkCopyOnWriteTable.java:113) at org.apache.hudi.table.HoodieSparkCopyOnWriteTable.upsert(HoodieSparkCopyOnWriteTable.java:97) at org.apache.hudi.client.SparkRDDWriteClient.upsert(SparkRDDWriteClient.java:157) at org.apache.hudi.DataSourceUtils.doWriteOperation(DataSourceUtils.java:213) at org.apache.hudi.HoodieSparkSqlWriter$.write(HoodieSparkSqlWriter.scala:304) at org.apache.hudi.DefaultSource.createRelation(DefaultSource.scala:163) at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:45) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:75) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:73) at org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:84) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:115) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:232) at org.apache.spark.sql.execution.SQLExecution$.executeQuery$1(SQLExecution.scala:110) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:135) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:232) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:135) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:253) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:134) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:68) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:112) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:108) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:519) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:83) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:519) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:495) at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:108) at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:95) at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:93) at org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:136) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:848) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:382) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:355) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:239) at com.licious.dataplatform.datalake.pipelines.gold.MonthABTest$.execute(MonthABTest.scala:95) at com.licious.dataplatform.datalake.pipelines.gold.MonthABTest$.main(MonthABTest.scala:104) at com.licious.dataplatform.datalake.pipelines.gold.MonthABTest.main(MonthABTest.scala) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$2.run(ApplicationMaster.scala:740) Caused by: org.apache.spark.SparkException: 阶段失败导致作业终止: 阶段42.0中的任务0失败4次,最近一次失败: 丢失阶段42.0中的任务0.3 (TID 1214) (ip-10-1-5-17.ap-south-1.compute.internal 执行器8): org.apache.hudi.exception.HoodieIOException: 无法读取Parquet文件s3://ls-dp/prod/datalake/data/test/customer/monthly_user_properties/pt=2024-01/8df75d2d-023f-48a0-b372-af726deb6c41-0_0-62-2170_20240102172043917.parquet at org.apache.hudi.common.util.ParquetUtils.getHoodieKeyIterator(ParquetUtils.java:181) at org.apache.hudi.common.util.ParquetUtils.fetchHoodieKeys(ParquetUtils.java:196) at org.apache.hudi.common.util.ParquetUtils.fetchHoodieKeys(ParquetUtils.java:147) at org.apache.hudi.io.HoodieKeyLocationFetchHandle.locations(HoodieKeyLocationFetchHandle.java:62) at org.apache.hudi.index.simple.HoodieSimpleIndex.lambda$fetchRecordLocations$33972fb4$1(HoodieSimpleIndex.java:155) at org.apache.hudi.data.HoodieJavaRDD.lambda$flatMap$a6598fcb$1(HoodieJavaRDD.java:117) at org.apache.spark.api.java.JavaRDDLike.$anonfun$flatMap$1(JavaRDDLike.scala:125) at scala.collection.Iterator$$anon$11.nextCur(Iterator.scala:486) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:492) at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460) at org.apache.spark.shuffle.sort.UnsafeShuffleWriter.write(UnsafeShuffleWriter.java:183) at org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52) at org.apache.spark.scheduler.Task.run(Task.scala:133) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:506) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1474) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:509) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) Caused by: java.io.FileNotFoundException: 找不到文件或目录's3://ls-dp/prod/datalake/data/test/customer/monthly_user_properties/pt=2024-01/8df75d2d-023f-48a0-b372-af726deb6c41-0_0-62-2170_20240102172043917.parquet' at com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem.getFileStatus(S3NativeFileSystem.java:521) at com.amazon.ws.emr.hadoop.fs.EmrFileSystem.getFileStatus(EmrFileSystem.java:613) at org.apache.parquet.hadoop.ParquetReader$Builder.build(ParquetReader.java:337) at org.apache.hudi.common.util.ParquetUtils.getHoodieKeyIterator(ParquetUtils.java:178) ... 20 more
解决方案建议
- 清理Hudi元数据中对已删除Parquet文件的引用:
- 检查
.hoodie/commits、.hoodie/compaction、.hoodie/rollbacks目录下的所有文件,删除包含pt=2024-01分区相关文件路径的记录。 - 若使用Simple Index,需手动删除
.hoodie/index目录中对应分区的索引文件;若使用Bloom Index,清理对应分区的Bloom过滤文件。
- 检查
- 执行Hudi的
REPAIR TABLE命令,强制Hudi重新同步元数据与实际文件系统状态:spark.sql("REPAIR TABLE your_hudi_table_name") - 若上述操作无效,可尝试直接重新创建该分区的数据:
- 确保写入操作使用
upsert模式,并且设置hoodie.datasource.write.ignore.existing.partitions为true(Hudi 0.10+版本支持),跳过对现有分区元数据的校验。 - 重新写入
pt=2024-01分区的数据,让Hudi重新生成该分区的元数据和数据文件。
- 确保写入操作使用
内容的提问来源于stack exchange,提问作者Priyanshu Sharma
相关产品推荐
相关产品推荐

