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

PySpark合并DataFrame并覆盖写入GCS路径时遇文件未找到错误

问题

开发时需要将脚本生成的DataFrame与Google Cloud Storage(GCS)存储桶中已有的DataFrame合并。当目标文件存在(count>0)时,将newDf和读取的df执行union并去重,再以overwrite模式写回原路径;文件不存在时直接写入正常,但写回原路径时报错。

代码示例

# newDf是通过脚本逻辑生成的DataFrame,count用于判断文件是否存在
if count > 0:
    df = spark.read.option('header', 'true').csv(path)
    dfUnion = df.union(newDf).dropDuplicates(["tempColumn"])
    dfUnion.write.mode("overwrite").option("header", "true").csv(path)
else:
    newDf.write.mode("overwrite").option("header", "true").csv(path)

报错信息

Traceback (most recent call last):
  File "/tmp/a78b22f8-4000-4a0a-8e7a-0c4e65e748d7/_init_.py", line 17, in <module>
    run()
  File "/tmp/a78b22f8-4000-4a0a-8e7a-0c4e65e748d7/_init_.py", line 13, in run
    script.run(utilities)
  File "/tmp/a78b22f8-4000-4a0a-8e7a-0c4e65e748d7/script.py", line 213, in run
    dfUnion.write.mode("overwrite").option("header", "true").csv(path)
  File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 1372, in csv
  File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py", line 1304, in __call__
  File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 111, in deco
    
  File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py", line 326, in get_return_value
py4j.protocol.Py4JJavaError: 调用o32819.csv时发生错误。
: org.apache.spark.SparkException: 作业已中止。
    at org.apache.spark.sql.execution.datasources.FileFormatWriter$.write(FileFormatWriter.scala:231)
    ...
Caused by: org.apache.spark.SparkException: 由于阶段失败导致作业中止:阶段1674.0中的任务2失败4次,最近一次失败:丢失阶段1674.0中的任务2.3(TID 2577)(job-w-1执行器7):java.io.FileNotFoundException: 未找到项目: 'gs://bkt-name/path/part-00000-6d328b30-ea18-4eef-8148-d4d732c84893-c000.csv'。注意,当前版本可能仍可用,但请求的生成版本已被删除。
可能底层文件已更新。您可以通过在SQL中运行'REFRESH TABLE tableName'命令或重新创建相关的Dataset/DataFrame来显式清除Spark缓存。
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.org$apache$spark$sql$execution$datasources$FileScanRDD$$anon$$readCurrentFile(FileScanRDD.scala:124)
    ...
驱动程序栈跟踪:
    at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2304)
    ...
Caused by: java.io.FileNotFoundException: 未找到项目: 'gs://bkt-name/path/part-00000-6d328b30-ea18-4eef-8148-d4d732c84893-c000.csv'。注意,当前版本可能仍可用,但请求的生成版本已被删除。
可能底层文件已更新。您可以通过在SQL中运行'REFRESH TABLE tableName'命令或重新创建相关的Dataset/DataFrame来显式清除Spark缓存。
    at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.org$apache$spark$sql$execution$datasources$FileScanRDD$$anon$$readCurrentFile(FileScanRDD.scala:124)
    ...

解决方案

错误核心原因:Spark执行overwrite写入时会先删除目标路径下的所有文件,但此时之前读取的df对应的RDD仍在尝试读取已被删除的旧文件,导致文件找不到。

方法1:缓存合并后的DataFrame再写入

先将合并后的DataFrame缓存到内存,确保数据已加载完成后再执行写入:

if count > 0:
    df = spark.read.option('header', 'true').csv(path)
    dfUnion = df.union(newDf).dropDuplicates(["tempColumn"])
    # 缓存并触发数据加载到内存
    dfUnion.cache()
    dfUnion.count()
    # 执行写入
    dfUnion.write.mode("overwrite").option("header", "true").csv(path)
    # 释放缓存
    dfUnion.unpersist()
else:
    newDf.write.mode("overwrite").option("header", "true").csv(path)

方法2:写入临时路径再替换原路径

先将合并后的数据写入临时路径,再替换原路径的文件:

if count > 0:
    temp_path = f"{path}_temp"
    df = spark.read.option('header', 'true').csv(path)
    dfUnion = df.union(newDf).dropDuplicates(["tempColumn"])
    
    # 写入临时路径
    dfUnion.write.mode("overwrite").option("header", "true").csv(temp_path)
    
    # 操作HDFS文件系统替换路径
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    fs = spark.sparkContext._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf)
    src_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(temp_path)
    dest_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(path)
    
    # 删除原路径,移动临时路径文件到原路径
    fs.delete(dest_path, True)
    fs.rename(src_path, dest_path)
else:
    newDf.write.mode("overwrite").option("header", "true").csv(path)

方法3:使用Spark SQL临时表完成合并

通过临时表的方式避免直接依赖文件路径的读取逻辑:

if count > 0:
    df = spark.read.option('header', 'true').csv(path)
    df.createOrReplaceTempView("existing_data")
    newDf.createOrReplaceTempView("new_data")
    
    # 用SQL完成合并去重
    dfUnion = spark.sql("""
        SELECT DISTINCT * FROM (
            SELECT * FROM existing_data
            UNION ALL
            SELECT * FROM new_data
        ) t
    """)
    dfUnion.write.mode("overwrite").option("header", "true").csv(path)
else:
    newDf.write.mode("overwrite").option("header", "true").csv(path)

内容的提问来源于stack exchange,提问作者QueryQuasar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 04:40:44