Spark DataFrame写入S3报文件未找到,关联Hive分区表覆盖操作
我有一个按date和site列分区的Hive表,用它创建DataFrame后对前一日数据做计算,然后覆盖原表的操作是成功的。但把最终的DataFrame写入S3桶时,出现了文件未找到的错误,报错指向的是已经被覆盖的前一日文件。如果我先把DataFrame写入S3再覆盖Hive表,就能正常运行。想请教下为什么写入S3会关联到Hive表的分区文件?
报错信息
java.io.FileNotFoundException: No such file or directory: s3://bucket_1/DM/web_fact_tbl/local_dt=2018-05-10/site_name=ABC/part-00000-882a6e29-eb6a-477c-8b88-6fe853956674.c000
相关代码
fact_tbl = spark.table('db.web_fact_tbl') fact_lkp = fact_tbl.filter(fact_tbl['local_dt']=='2018-05-10') fact_join = fact_lkp.alias('a').join(fact_tbl.alias('b'),(col('a.id') == col('b.id')),"inner").select('a.*') fact_final = fact_join.union(fact_tbl) fact_final.coalesce(2).createOrReplaceTempView('cwf') spark.sql('INSERT OVERWRITE TABLE dm.web_fact_tbl PARTITION (local_dt, site_name) \ SELECT * FROM cwf') fact_final.write.csv('s3://bucket_1/yahoo')
这事儿其实和Spark的惰性求值机制以及DataFrame的血统(lineage)直接相关,我给你拆解清楚:
DataFrame的血统没断开:你创建的
fact_final是基于原始Hive表fact_tbl做join、union得到的,Spark不会立刻执行计算,而是把这些操作记录成一个执行计划(也就是血统)。当你执行INSERT OVERWRITE覆盖Hive表后,Hive分区对应的S3文件已经被删除了,但fact_final的执行计划还依赖原始的Hive表文件——等到你调用fact_final.write.csv时,Spark才会去执行这个血统里的操作,这时候它找不到已经被覆盖删除的旧文件,自然就报了错。先写S3再覆盖Hive表为啥没问题:因为先执行
fact_final.write.csv时,Spark会先触发计算,把fact_final的数据计算出来并写入S3,这时候已经把所有需要的原始文件都读取过了,之后再覆盖Hive表,就不会影响已经完成的计算了。
给你两个可行的解决办法:
提前触发计算,断开血统依赖
在覆盖Hive表之前,先把fact_final的数据物化(比如用cache()+count()强制触发计算),这样Spark会提前把数据计算出来存在内存/磁盘里,后续操作就不再依赖原始Hive表的文件了:fact_tbl = spark.table('db.web_fact_tbl') fact_lkp = fact_tbl.filter(fact_tbl['local_dt']=='2018-05-10') fact_join = fact_lkp.alias('a').join(fact_tbl.alias('b'),(col('a.id') == col('b.id')),"inner").select('a.*') fact_final = fact_join.union(fact_tbl).cache() # 触发计算,把数据物化 fact_final.count() # 之后再执行覆盖和写入S3 fact_final.coalesce(2).createOrReplaceTempView('cwf') spark.sql('INSERT OVERWRITE TABLE dm.web_fact_tbl PARTITION (local_dt, site_name) \ SELECT * FROM cwf') fact_final.write.csv('s3://bucket_1/yahoo')调整执行顺序,先写S3再覆盖表
既然这个顺序能正常运行,直接把写入S3的操作放在覆盖Hive表之前就行,这样计算完成后再覆盖表,就不会出现文件找不到的问题:fact_tbl = spark.table('db.web_fact_tbl') fact_lkp = fact_tbl.filter(fact_tbl['local_dt']=='2018-05-10') fact_join = fact_lkp.alias('a').join(fact_tbl.alias('b'),(col('a.id') == col('b.id')),"inner").select('a.*') fact_final = fact_join.union(fact_tbl) # 先写入S3,触发计算 fact_final.write.csv('s3://bucket_1/yahoo') # 再覆盖Hive表 fact_final.coalesce(2).createOrReplaceTempView('cwf') spark.sql('INSERT OVERWRITE TABLE dm.web_fact_tbl PARTITION (local_dt, site_name) \ SELECT * FROM cwf')
另外提个小优化点:coalesce(2)是窄依赖操作,建议放在最后一步执行,避免影响之前步骤的并行度,但这不是导致报错的原因~
内容的提问来源于stack exchange,提问作者yahoo

