AWS Pyspark Glue作业中Spark overwrite模式触发文件不存在错误如何解决
问题根因
你遇到的错误由Spark懒执行机制导致:所有转换算子不会立即执行,直到调用finalDF.write这类action算子时才会触发整条计算链路的执行。overwrite模式写Parquet的第一步是清空目标目录transformedPath,等后续执行到spark.read.parquet(transformedPath)的读取逻辑时,该目录已经被删除,因此抛出文件不存在错误。
最优解决方案(零路径改造、最低改造成本)
仅需在增量分支的聚合逻辑后、写操作前,对聚合结果加缓存并触发一次立即计算,让数据暂存在Glue作业的内存/本地磁盘中,不再依赖S3上的源文件即可解决问题,修改后代码如下:
elif (incrementalLoad == str(1)): app8.write.mode("append").parquet(transformedPath)#loc1 print(":::Incremental Transformed data has been written::::::::") transformedData = spark.read.parquet(transformedPath) print("::::::Transformed data has been written:::::") finalDF = transformedData.groupBy(col("mobilenumber")).agg( sum(col("times_app_uninstalled")).alias("times_app_uninstalled"), sum(col("times_uninstall_l30d")).alias("times_uninstall_l30d"), sum(col("times_uninstall_l60d")).alias("times_uninstall_l60d"), sum(col("times_uninstall_l90d")).alias("times_uninstall_l90d"), sum(col("times_uninstall_l180d")).alias("times_uninstall_l180d"), sum(col("times_uninstall_l270d")).alias("times_uninstall_l270d"), sum(col("times_uninstall_l365d")).alias("times_uninstall_l365d"), max(col("latest_uninstall_date")).alias("latest_uninstall_date"), min(col("first_uninstall_date")).alias("first_uninstall_date")) # 新增两行代码即可解决问题 finalDF = finalDF.cache() # 触发action提前完成全量聚合计算,数据存入缓存,不再依赖S3源文件 finalDF.count() finalDF.write.mode("overwrite").parquet(transformedPath)#loc1
如果聚合后数据量较大内存存不下,可以将cache()替换为persist(StorageLevel.DISK_ONLY),改用Glue作业本地磁盘暂存数据,效果完全一致。
进阶性能优化方案(适合后续数据量持续增长场景)
如果后续transformed路径下的明细数据量持续增大,全量读取明细做聚合的效率会不断降低,可以改造为增量聚合逻辑,进一步降低计算开销:
- 全量加载阶段:保持原有逻辑,直接将聚合后的最终结果写入
transformedPath - 增量加载阶段:
- 先读取
transformedPath的历史聚合结果存为historyAggDF - 对增量转换结果
app8先按手机号做一次预聚合,得到incrAggDF - 将
historyAggDF与incrAggDF做union,再按手机号做二次聚合得到最新全量聚合结果 - 按上述方案缓存后overwrite写入
transformedPath
该方案仅需读取历史聚合结果和增量数据,无需读取全量明细,性能会有明显提升。
- 先读取
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

