S3中同时覆盖与读取操作引发FileNotFoundException问题排查
你遇到的FileNotFoundException本质是Spark文件元数据缓存和S3最终一致性的叠加问题,再加上Spark的惰性执行特性:
- 第一次读取
sourcePath时,Spark会缓存该路径的文件元数据(即使是惰性执行,部分元数据可能提前拉取) - 手动删除
sourcePath并将临时目录重命名为sourcePath的操作,S3需要一定时间同步元数据(最终一致性) - 后续读取同一
sourcePath时,Spark可能仍使用缓存的旧元数据,或者S3还没同步完成,导致找不到文件
之前尝试的thread.sleep(5000)无法适配S3元数据同步的波动时间,spark.catalog.clearCache()只针对Hive表缓存,对直接读取路径的场景无效。
1. 强制刷新路径元数据
在读取sourcePath之前,执行路径刷新命令,让Spark重新拉取最新的文件列表:
// Scala 示例 spark.sql(s"REFRESH PATH '$sourcePath'")
# Python 示例 spark.sql(f"REFRESH PATH '{sourcePath}'")
这个命令会直接触发Spark去S3重新扫描指定路径的文件,覆盖缓存的旧元数据。
2. 完全重新初始化读取逻辑
确保每次读取sourcePath时都创建全新的DataFrame,不要复用之前的DataFrame/Dataset对象:
// 错误:复用之前的读取逻辑(可能携带旧元数据) // val df = spark.read.json(sourcePath) // ... 处理后再次使用 df // 正确:每次读取都重新初始化 val freshDf = spark.read.json(sourcePath) // 后续处理使用 freshDf
Spark的惰性执行会让旧DataFrame保留最初的元数据引用,重新创建对象能彻底避免这个问题。
3. 优化S3操作的原子性
替换手动删除+重命名的逻辑,利用Spark的原子写入特性,或者S3的原子操作:
方案A:直接覆盖写入
sourcePath
如果你需要将多文件合并为单文件,直接用Overwrite模式写入目标路径,Spark会自动处理临时文件和原子替换,比手动操作更可靠:spark.read .json(sourcePath) .coalesce(1) .write .mode(SaveMode.Overwrite) .json(sourcePath)注意:这种方式会直接覆盖原路径下的所有文件,确保业务逻辑允许。
方案B:S3原子重命名
使用AWS SDK执行原子重命名(S3没有原生重命名,实际是复制+删除,但可以通过SDK确保操作的原子性),避免中间状态暴露给Spark:import boto3 s3 = boto3.resource('s3') bucket = s3.Bucket('your-bucket-name') # 复制临时目录下的文件到目标路径 for obj in bucket.objects.filter(Prefix=tempTarget1.replace('s3://your-bucket-name/', '')): dest_key = obj.key.replace(tempTarget1.split('/')[-2], sourcePath.split('/')[-2]) bucket.copy({'Bucket': 'your-bucket-name', 'Key': obj.key}, dest_key) obj.delete() # 删除原sourcePath下的旧文件 for obj in bucket.objects.filter(Prefix=sourcePath.replace('s3://your-bucket-name/', '')): obj.delete()
4. 禁用Spark文件元数据缓存
如果你的作业频繁变更同一路径的内容,可以全局禁用Spark的文件元数据缓存,代价是每次读取都要扫描路径,性能略有下降:
// 在SparkSession初始化时设置 val spark = SparkSession.builder() .config("spark.sql.files.cacheCatalogMetadata", "false") .getOrCreate()
这个配置会让Spark每次读取路径时都实时扫描文件,不会缓存元数据。
5. 等待S3元数据同步完成
用循环检查替代固定sleep,直到S3返回目标路径存在文件再继续读取:
import boto3 import time def wait_for_s3_path(s3_path): s3 = boto3.client('s3') bucket, prefix = s3_path.replace('s3://', '').split('/', 1) while True: resp = s3.list_objects_v2(Bucket=bucket, Prefix=prefix) if 'Contents' in resp and len(resp['Contents']) > 0: break time.sleep(1) # 在重命名完成后调用 wait_for_s3_path(sourcePath) # 再执行读取操作 df = spark.read.json(sourcePath)
内容的提问来源于stack exchange,提问作者Ashit_Kumar

