修改Parquet Schema后覆盖写入S3遇FileNotFoundException求助
问题分析与解决方案
错误核心是Spark在overwrite模式写入Parquet到S3时,找不到指定文件s3://s3-tdh-infolake-test/vai/product_v1/Merged_20221027093416_product_v1_21,本质是S3作为对象存储的特性与Spark overwrite机制冲突:Spark会先删除目标路径下所有文件再写入新文件,但S3的最终一致性可能导致删除后仍有任务尝试读取已删除文件,或是路径依赖混淆引发读取错误。
解决步骤
改用临时路径过渡,避免直接覆盖
先将处理后的DataFrame写入临时路径,确保写入成功后再替换原路径:# 写入临时路径 df2.write.parquet('s3://s3-tdh-infolake-test/vai/product_v1/temp_Merged_product_v1_21')删除原目标路径
确认临时路径写入完成后,删除原目标路径(可通过Spark API操作):# Spark API方式删除原路径 from pyspark.sql import SparkSession spark = SparkSession.getActiveSession() hadoop_conf = spark._jsc.hadoopConfiguration() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) target_path = spark._jvm.org.apache.hadoop.fs.Path('s3://s3-tdh-infolake-test/vai/product_v1/Merged_202_product_v1_21') if fs.exists(target_path): fs.delete(target_path, recursive=True)重命名临时路径到目标路径
将临时路径的文件移动到原目标路径:temp_path = spark._jvm.org.apache.hadoop.fs.Path('s3://s3-tdh-infolake-test/vai/product_v1/temp_Merged_product_v1_21') fs.rename(temp_path, target_path)
额外检查点
- 核对路径拼写:代码中目标路径是
Merged_202_product_v1_21,但错误中提到的是Merged_20221027093416_product_v1_21,排查是否存在原文件路径与目标路径混淆的情况。 - 验证S3权限:确保Spark执行角色拥有目标路径的读写权限,权限不足也可能伪装成文件未找到的错误。
- 调整Spark协议:尝试设置
spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem,使用S3A协议替代旧的S3N协议,提升兼容性。
内容的提问来源于stack exchange,提问作者Vaibhav Tyagi
相关产品推荐
相关产品推荐

