PySpark中Parquet文件转换后覆盖写入触发FileNotFound Error问题排查及正确写入方式咨询
问题分析与解决方案
为什么会出现FileNotFoundError?
你遇到的问题核心在于Spark同时读取和写入同一个本地路径时的资源冲突:
当你执行c.write.mode('overwrite').parquet('./a')时,Spark的overwrite模式会先删除目标路径./a下的所有文件,但此时你的DataFrame c是基于从./a读取的a和b生成的——任务执行过程中,部分计算任务还需要读取./a下的原始文件,可这些文件已经被提前删除了,自然就抛出了FileNotFoundError。
简单来说:读操作还没完成,写操作就把源文件删了,任务找不到依赖的文件就报错了。
解决覆盖写入的问题
解决这个问题的核心思路是避免读写同一路径的冲突,推荐两种方案:
方案1:使用临时路径中转
先把转换后的结果写入临时路径,待写入完成后再替换原路径:
from pyspark.sql import SparkSession from pyspark.sql.functions import * import pandas as pd import shutil import os spark = SparkSession.builder.appName('test').getOrCreate() # 创建测试数据 a = spark.createDataFrame(pd.DataFrame({ 'x': [1, 2, 3], 'y':[1, 2, 3] })) a.write.mode('overwrite').parquet('./a') # 转换操作 a_df = spark.read.parquet('./a') b_df = spark.read.parquet('./a') c_df = a_df.union(b_df).withColumn('id', monotonically_increasing_id()) c_df.show() # 第一步:写入临时路径 temp_path = './temp_a' c_df.write.mode('overwrite').parquet(temp_path) # 第二步:替换原路径 # 先删除原路径 if os.path.exists('./a'): shutil.rmtree('./a') # 把临时路径重命名为原路径 shutil.move(temp_path, './a')
如果是在分布式环境(比如HDFS),可以用Hadoop命令来替换文件操作:
hadoop fs -rm -r ./a hadoop fs -mv ./temp_a ./a
方案2:提前将数据加载到内存(仅适合小数据量)
如果你的数据量很小,可以将读取的DataFrame缓存到内存中,这样后续的转换操作就不需要再读取磁盘文件了:
a_df = spark.read.parquet('./a').cache() # 缓存到内存 b_df = spark.read.parquet('./a').cache() c_df = a_df.union(b_df).withColumn('id', monotonically_increasing_id()) # 先触发缓存(执行一次action操作) c_df.count() # 再执行覆盖写入 c_df.write.mode('overwrite').parquet('./a')
⚠️ 注意:这个方法只适合小数据集,大数据量缓存会占用过多内存,导致性能问题。
Parquet正确的追加记录方式
如果只是要给现有Parquet文件追加新记录,直接使用mode('append')即可,但需要保证追加的DataFrame和原文件的Schema兼容:
- 字段名、数据类型要匹配
- 新增字段需设置为可空(nullable=true)
示例代码:
# 创建要追加的新数据 new_data = spark.createDataFrame(pd.DataFrame({ 'x': [4, 5], 'y': [4, 5] })) # 追加到原Parquet文件 new_data.write.mode('append').parquet('./a')
如果需要修改部分数据行,推荐的流程是:
- 读取原Parquet数据
- 过滤出需要修改的行,执行更新操作
- 将更新后的行与未修改的行合并
- 写入临时路径,再替换原路径(避免读写冲突)
比如修改x=2的行的y值:
original_df = spark.read.parquet('./a') # 未修改的行 unmodified_df = original_df.filter(original_df.x != 2) # 修改后的行 modified_df = original_df.filter(original_df.x == 2).withColumn('y', lit(100)) # 合并结果 final_df = unmodified_df.union(modified_df) # 写入临时路径后替换原路径 final_df.write.mode('overwrite').parquet('./temp_a') # 替换操作(和前面的方案1一致)
内容的提问来源于stack exchange,提问作者Zézouille
相关产品推荐
相关产品推荐

