Spark执行union转换后写DataFrame到CSV报文件不存在错误如何解决
问题根因
报错和表缓存、元数据异常没有关系,核心原因是读写同一路径+Spark懒执行机制+Overwrite写入模式的冲突:
- 你的代码先从
data/temporaryBasis路径读取数据生成dataTemporary,后续union逻辑依赖这个数据源 - 所有
read、drop、union、withColumn都是转换算子,仅构建执行计划,不会真正触发数据读取,所以你提前调用show能正常出结果——这时候原路径文件还没被改动 - 当你触发
write动作时,SaveMode.Overwrite模式会首先清空目标写入路径data/temporaryBasis下的所有原文件,再执行计算逻辑。这时候执行计划里需要读原路径下的文件做union,文件已经被删了,自然抛出FileNotFoundException。
不加union逻辑时,整个计算流程不依赖data/temporaryBasis路径的数据源,所以写入不会报错。
修复方法
选任意一种即可,生产环境优先选方案2稳定性最高:
- 方案1:预物化依赖同路径的结果数据,切断和原路径的关联
在写入前,把union后的最终结果先持久化到内存/磁盘,让后续计算不再依赖原路径的文件,示例代码:import org.apache.spark.sql.functions._ // 读取逻辑和之前一致 var data= spark.read.option("header","true").option("inferSchema","true").csv("data/dataset/mammography_id.csv") .drop("ID") var dataTemporary = spark.read.option("header","true").option("inferSchema","true").csv("data/temporaryBasis") .drop("ID") for(d<-dataTemporary.columns) if(d.contains("_bin")) dataTemporary=dataTemporary.drop(d) // 构建完最终结果后,先cache物化 val finalData = dataTemporary.union(data) .withColumn("ID",monotonically_increasing_id()) .withColumn("features", stringify(col("features"))) .persist() // 调用动作算子触发真正的计算,把数据存到缓存,不再依赖原路径文件 finalData.count() // 此时再写入同路径就不会冲突 finalData.write .mode(SaveMode.Overwrite) .option("header","true") .csv("data/temporaryBasis") // 写完释放缓存 finalData.unpersist() - 方案2:写临时路径后替换原目录(生产推荐)
完全避开读写同路径的问题,把结果先写到独立的临时路径,写入完成后通过文件系统操作删除原目录,再把临时目录重命名为目标路径即可,不存在缓存丢失、执行计划依赖的问题,稳定性最高。 - 方案3:小数据量场景可先把读取的路径数据收集到内存
如果data/temporaryBasis的数据量很小,可以在读取后直接调用collect()拉到Driver端,再重新转成DataFrame做后续计算,相当于把数据和原路径的绑定断开,不过数据量大会直接OOM,不推荐大数据量场景用。
之前操作的误区说明
- 报错提示的
REFRESH TABLE仅针对已经注册到Hive元数据的托管表/外部表场景,你当前是直接读写文件路径,没有把数据注册成表,刷表、手动建表的操作和当前问题无关,解决不了问题。 - 你用的
spark.catalog.createTable("newTable", "data/temporaryBasis")写法本身有误:createTable如果直接传路径,需要额外指定数据源格式、表结构、分区规则等参数,否则必然报错。如果需要基于DataFrame创建表,直接用df.write.saveAsTable("newTable")即可,不需要手动调用createTable。
内容的提问来源于stack exchange,提问作者informatic_aa
相关产品推荐
相关产品推荐

