Pyspark Glue作业写Parquet文件时S3存储桶被删除问题咨询
问题描述
我开发了一个Pyspark Glue作业用于加载全量/增量数据集,运行正常。加载数据集后我需要执行若干aggregations操作,再以*"overwrite"/"append"*模式写入到指定位置。为此我编写了如下代码:
maxDateValuePath = "s3://...../maxValue/" outputPath = "s3://..../complete-load/" aggregatedPath = "s3://...../aggregated-output/" fullLoad = "" aggregatedView = "" completeAggregatedPath = "s3://...../aggregated-output/step=complete-load/" incrAggregatedPath = "s3://....../aggregated-output/step=incremental-load/" aggregatedView="" data.createOrReplaceTempView("data") aggregatedView = spark.sql(""" select catid,count(*) as number_of_catids from data group by catid""") if (incrementalLoad == str(0)): aggregatedView = aggregatedView.withColumn("created_at", current_timestamp()) aggregatedView.write.mode("overwrite").parquet(completeAggregatedPath) elif (incrementalLoad == str(1)): aggregatedView = aggregatedView.withColumn("created_at", current_timestamp()) log.info("step 123: " + str(aggregatedView.count())) aggregatedView.write.mode("append").parquet(completeAggregatedPath) aggregatedView = spark.read.parquet(completeAggregatedPath) log.info("step 126: " + str(aggregatedView.count())) w = Window.partitionBy("catid").orderBy(col("created_at").desc()) aggregatedView = aggregatedView.withColumn("rw", row_number().over(w)).filter(col("rw") == lit(1)).drop( "rw") log.info("step 130: " + str(aggregatedView.count())) log.info(aggregatedView.orderBy(col("created_at").desc()).show()) print("::::::::::::before writing::::::::::::::") aggregatedView.write.mode("overwrite").parquet(incrAggregatedPath)
其中0和1分别代表全量加载/增量加载。我在写入转换后的数据集前新增了created_at字段,用于处理增量数据集写入后的最新聚合记录,避免产生重复数据。
所有功能运行均符合预期,但问题出现在增量流程中执行aggregatedView.write.mode("overwrite").parquet(aggregatedPath)代码写入数据集时,误以为S3存储桶被删除,操作抛出如下*error*:
Caused by: java.io.FileNotFoundException: File not present on S3 It is possible the underlying files have been updated. You can explicitly invalidate the cache in Spark by running 'REFRESH TABLE tableName' command in SQL or by recreating the Dataset/DataFrame involved.
问题解答
首先明确:S3存储桶本身没有被删除,报错是Spark读取不到指定路径下的文件,并非存储桶被移除。
触发问题的核心原因有两个:
- Overwrite模式的默认删除逻辑
aggregatedPath是completeAggregatedPath和incrAggregatedPath的父目录,Spark用overwrite模式写入路径时,默认会先清空目标路径下的所有内容再写入新文件。你执行写入aggregatedPath的操作时,会先删除该路径下已存在的step=complete-load/、step=incremental-load/两个子目录的全部内容。
而Spark的DataFrame是惰性计算的,你之前读取completeAggregatedPath生成的aggregatedView在执行写入动作时才会真正触发计算,此时completeAggregatedPath下的文件已经被overwrite操作提前删除,自然会报文件不存在的错误。 - Spark文件元数据缓存
Spark默认会缓存读取过的路径的文件列表元数据,你在append写入completeAggregatedPath后立刻读取该路径,可能拿到的是旧的缓存元数据,后续文件被删除后执行实际计算时,就会找不到对应的实际文件。
修复方案
- 不要直接写入父路径
aggregatedPath,你已经拆分了全量、增量的独立子路径,直接写入对应子路径incrAggregatedPath即可,不会影响其他目录的内容。 - 读取
completeAggregatedPath后对DataFrame执行cache()操作,避免后续计算重复读取已经被修改的路径:
aggregatedView = spark.read.parquet(completeAggregatedPath).cache()
- 读取路径前可以主动刷新元数据,避免使用旧的缓存信息:
spark.catalog.refreshByPath(completeAggregatedPath) aggregatedView = spark.read.parquet(completeAggregatedPath)
- 如果确实需要覆盖写入分区路径,可以配置动态分区覆盖参数,避免删除全路径内容:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "DYNAMIC")
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

