You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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')

如果需要修改部分数据行,推荐的流程是:

  1. 读取原Parquet数据
  2. 过滤出需要修改的行,执行更新操作
  3. 将更新后的行与未修改的行合并
  4. 写入临时路径,再替换原路径(避免读写冲突)

比如修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 17:47:47