AWS EMR批处理中Spark persist后write触发FileNotFoundException排查
AWS EMR Spark批处理任务FileNotFoundException排查方案
可能的原因及解决办法
1. 读写同一路径引发的重计算冲突
如果你的updateDF.write目标路径和spark.read.load(paths)的源路径是同一个,这是最可能的触发原因:
- Spark的
overwrite模式会在写入过程中清理或替换目标路径的旧文件;而updateDF.persist()的数据如果因为内存不足被驱逐,Spark会自动重新从源路径读取数据进行计算,此时源路径的文件已经被修改/删除,就会抛出文件找不到的异常。 - 解决办法:
- 严格分离读写路径:先将处理后的数据写入临时S3路径,任务完成后通过AWS CLI命令
aws s3 mv s3://temp-path/ s3://target-path/ --recursive替换原路径。 - 提升persist存储级别:使用更稳定的存储级别避免数据被驱逐,比如:
from pyspark import StorageLevel updateDF.persist(StorageLevel.MEMORY_AND_DISK_SER)
- 严格分离读写路径:先将处理后的数据写入临时S3路径,任务完成后通过AWS CLI命令
2. S3最终一致性特性导致的文件不可见
S3作为对象存储,默认是最终一致性:文件写入/覆盖后,部分区域或节点可能存在同步延迟,Spark executor所在的节点可能暂时无法读取到新生成的文件。
- 解决办法:
- 开启S3强一致性:在EMR集群配置中启用
EMRFS Consistent View,或者添加Spark配置spark.hadoop.fs.s3a.consistent=true。 - 写入完成后添加短暂等待:如果后续操作依赖写入后的文件,可添加
time.sleep(30)(根据实际情况调整时长),给S3足够的同步时间。
- 开启S3强一致性:在EMR集群配置中启用
3. EMRFS元数据缓存过期
EMRFS默认会缓存S3的文件元数据,当S3文件被修改后,缓存的元数据没有及时更新,导致executor尝试读取已不存在的旧文件。
- 解决办法:
- 禁用EMRFS元数据缓存:在Spark配置中添加
spark.hadoop.fs.s3a.metadatacache.impl=org.apache.hadoop.fs.s3a.NullMetadataCache。 - 手动刷新缓存:在读写操作之间执行
emrfs sync s3://{bucket_name}/{folder_name}/命令(可通过EMR Step或代码调用subprocess执行)。
- 禁用EMRFS元数据缓存:在Spark配置中添加
4. 任务重试时的临时文件丢失
当Spark task失败重试时,S3上的临时part文件可能因为S3生命周期规则、EMR临时文件清理机制被删除,导致重试任务找不到文件。
- 解决办法:
- 检查S3生命周期规则:确保任务执行期间,目标路径的文件不会被自动删除。
- 调整Spark重试配置:降低
spark.task.maxFailures的值避免无意义重试;或配置Spark使用本地临时目录存储中间数据,减少对S3临时文件的依赖。
内容的提问来源于stack exchange,提问作者Jisu Choi
相关产品推荐
相关产品推荐

