Spark写入S3新路径抛出FileNotFoundException的问题排查与解决
Spark写入S3 Parquet时偶发FileNotFoundException问题排查与解决
问题场景
我们需要将DataFrame以Parquet格式、overwrite模式写入S3的全新文件夹,使用的代码如下:
df.write .option("dateFormat", "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'") .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'") .option("maxRecordsPerFile", maxRecordsPerFile) .mode("overwrite") .format(format) .save(output)
但偶尔会抛出FileNotFoundException,完整报错栈如下:
Caused by: java.io.FileNotFoundException: No such file or directory 's3://BUCKET/snapshots/FOLDER/_bid_9223370368440344985/part-00020-693dfbcb-74e9-45b0-b892-0b19fa92365c-c000.snappy.parquet' It is possible the underlying files have been updated. You can explicitly invalidate the cache in Spark by running 'REFRESH TABLE tableName' command in SQL or by recreating the Dataset/DataFrame involved. at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.org$apache$spark$sql$execution$datasources$FileScanRDD$$anon$$readCurrentFile(FileScanRDD.scala:131) at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.nextIterator(FileScanRDD.scala:182) at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:109) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408) at scala.collection.Iterator$$anon$13.hasNext(Iterator.scala:461) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408) at org.apache.spark.sql.execution.aggregate.HashAggregateExec$$anonfun$doExecute$1$$anonfun$4.apply(HashAggregateExec.scala:104) at org.apache.spark.sql.execution.aggregate.HashAggregateExec$$anonfun$doExecute$1$$anonfun$4.apply(HashAggregateExec.scala:101) at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$26.apply(RDD.scala:853) at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$26.apply(RDD.scala:853) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:49) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324) at org.apache.spark.rdd.RDD.iterator(RDD.scala:288) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:49) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324) at org.apache.spark.rdd.RDD.iterator(RDD.scala:288) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:96) at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:53) at org.apache.spark.scheduler.Task.run(Task.scala:109) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345) 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)
针对该场景,以下是问题的解答:
1. 写入全新S3路径(无任何读取操作)时,为何会抛出该异常?
核心原因是S3的最终一致性特性加上Spark 2.3.x写入机制的冲突:
- Spark写入时会先将数据写入目标路径下的临时文件夹(如报错中的
_bid_xxxx),完成写入后再将临时文件移动到目标路径并删除临时文件夹。 - S3属于最终一致性存储,文件移动操作完成后,部分节点可能因S3的同步延迟或客户端缓存,无法立即看到移动后的文件。此时如果写入流程内部的元数据校验步骤尝试读取这些文件,就会触发
FileNotFoundException。 - 旧版
s3://客户端(EMR 5.x中的原生S3客户端)自带文件列表缓存,缓存未及时更新也会加剧这个问题。
2. 如何修复此问题?
针对无读写并发的场景,可采用以下方案:
- 改用原子性路径重命名逻辑:先写入独立的临时路径,确认写入完成后再将整个路径重命名到目标路径(S3的路径重命名是原子操作),替代Spark默认的临时文件移动逻辑。示例代码:
val tempOutput = output + "_temp_" + System.currentTimeMillis() df.write .option("dateFormat", "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'") .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'") .option("maxRecordsPerFile", maxRecordsPerFile) .mode("overwrite") .format(format) .save(tempOutput) // 重命名临时路径到目标路径 val fs = FileSystem.get(new URI(output), spark.sparkContext.hadoopConfiguration) fs.rename(new Path(tempOutput), new Path(output)) - 禁用S3客户端缓存:在Spark配置中添加以下参数,关闭文件系统的元数据缓存:
spark.hadoop.fs.s3a.metadatastore.impl org.apache.hadoop.fs.s3a.metadata.InMemoryMetadataStore spark.hadoop.fs.s3a.list.status.cache.enable false - 添加短暂等待(应急方案):在
save操作后添加1-2秒等待,给S3足够时间完成一致性同步,适合临时快速解决问题。
3. 关于环境(Spark 2.3.2、EMR-5.18.1、Scala)的适配说明
EMR 5.18.1自带的Spark 2.3.2对S3的支持存在已知的一致性和缓存问题,上述的路径重命名、禁用缓存、切换s3a客户端的方案在该环境下完全可行,EMR 5.x已内置s3a客户端的支持。
4. 是否需要切换为s3n或s3a?该操作是否有效?
- 不建议使用s3n:s3n是旧版S3客户端,性能、一致性支持均不如s3a,且已被官方废弃。
- 强烈建议切换为s3a:s3a是Hadoop官方推荐的S3客户端,相比原生
s3://客户端,它更好地适配了S3的最终一致性特性,支持更高效的文件操作,缓存机制也更可控。切换后配合禁用缓存的配置,能有效降低这类异常的发生概率。切换方式很简单:将输出路径前缀从s3://改为s3a://,同时添加上述s3a相关配置即可。
内容的提问来源于stack exchange,提问作者chendu
相关产品推荐
相关产品推荐

