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

Spark使用partitionBy写入Parquet时临时目录抛出FileAlreadyExistsException

Spark本地模式写入分区Parquet时临时目录触发FileAlreadyExistsException的原因及解决

问题场景

尝试将CSV文件转换为按CITY分区的Parquet数据集,执行如下代码:

master = "local[*]"
app_name = "convert_to_parquet"
spark = (
    SparkSession.builder
    .appName(app_name)
    .master(master)
    .getOrCreate()
)

csv_path = "<csv-path>"
in_df = spark.read.option("inferSchema", "true").option("header", "true").csv(csv_path)
out_df = in_df.selectExpr("trim(OBJECTID) AS ID",
                          "trim(NAME) AS NAME",
                          "trim(CITY) AS CITY",
                          "trim(STATE) AS STATE", "X", "Y")
out_path = "<out-dir>"
# 尝试提前删除输出目录
shutil.rmtree(out_path, ignore_errors=True)

out_df.write.partitionBy(["CITY"]).parquet(out_path)

运行后触发临时目录的文件已存在错误:

23/02/24 01:54:22 ERROR Utils: Aborting task                        (0 + 1) / 1]
org.apache.hadoop.fs.FileAlreadyExistsException: File already exists: file: <out-dir>/_temporary/0/_temporary/attempt_202302240154215533456461598462887_0002_m_000000_2/CITY=Apex/part-00000-0de20945-1012-4a7a-b183-c4235717a0a2.c000.snappy.parquet
        at org.apache.hadoop.fs.RawLocalFileSystem.create(RawLocalFileSystem.java:421)
        at org.apache.hadoop.fs.RawLocalFileSystem.create(RawLocalFileSystem.java:459)
        ...

原因分析

  1. 本地文件系统与HDFS的原子性差异:Spark写入逻辑依赖HDFS的原子重命名操作保证任务重试一致性,但本地文件系统(如ext4、NTFS)不支持该特性。当任务因IO延迟、调度波动重试时,之前生成的临时文件未被彻底清理,导致重复创建报错。
  2. 手动删除与Spark API的协同问题:shutil.rmtree是本地文件系统操作,而Spark通过Hadoop FS API处理文件,两者状态同步存在延迟。即便手动删除目录,Spark的缓存或未刷新的文件系统状态仍可能引发冲突。
  3. 本地多线程并行写入的锁缺失:local[*]模式启用多线程并行执行任务,多个线程同时写入同一分区目录时,本地文件系统无分布式锁协调,会出现多线程尝试创建同一临时文件的情况。

解决办法

  • 限制并行度为单线程:将Spark master设置为local[1],避免多线程同时操作本地分区目录:
    master = "local[1]"
    
  • 使用Spark Hadoop API删除目录:替换shutil.rmtree为Spark原生文件系统操作,确保与写入逻辑协同:
    # 获取Hadoop文件系统实例
    fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
    out_path_hadoop = spark._jvm.org.apache.hadoop.fs.Path(out_path)
    if fs.exists(out_path_hadoop):
        fs.delete(out_path_hadoop, True)  # True表示递归删除所有子目录
    
  • 启用覆盖写入模式:显式指定mode("overwrite"),让Spark自行处理目录清理与文件覆盖,这是最可靠的方式:
    out_df.write.mode("overwrite").partitionBy(["CITY"]).parquet(out_path)
    
  • 排查分区字段异常值:确认trim(CITY)后的字段值无特殊字符(如空格、斜杠),避免路径解析错误导致临时文件异常。

内容的提问来源于stack exchange,提问作者Bing-hsu Gao

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 22:48:20