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
相关产品推荐
相关产品推荐

