在Scala循环中执行.sql文件并将结果存入DataFrame的问题
解决Spark Shell中读取SQL文件并将结果存入DataFrame的问题
问题分析
你现有代码存在几个关键问题:
- 语法错误:
val sqlQuery -应为val sqlQuery = - 拆分SQL语句时未过滤空字符串,可能触发无效空查询
- 原代码仅调用
show()打印结果,未保存spark.sql()返回的DataFrame实例 - 使用
sc.textFile()时的错误是误传参数(collect()无需额外参数),且分布式读取本地小文件属于资源浪费
正确实现步骤
1. 正确读取并拆分SQL文件
使用scala.io.Source读取本地SQL文件,拆分语句并过滤无效内容:
import scala.io.Source // 读取SQL文件内容,拆分语句,过滤空语句和纯空格语句 val sqlStatements = Source.fromFile("metadata/usr_home/2_pre_checks.sql") .mkString .split(";") .map(_.trim) .filter(_.nonEmpty)
2. 执行SQL并获取DataFrame
将每个SQL语句的执行结果存入DataFrame列表,若结果结构一致可合并为单个DataFrame:
// 执行所有SQL,得到DataFrame列表 val resultDfs = sqlStatements.map(query => spark.sql(query)) // 若所有查询结果结构相同,合并为一个DataFrame(适配你的union all场景) val combinedDf = resultDfs.reduce(_ unionAll _) // 查看合并后的结果 combinedDf.show(false)
3. 保存DataFrame
可将结果写入Hive表或文件系统:
// 写入Hive表 combinedDf.write.mode("overwrite").saveAsTable("your_result_table") // 写入Parquet文件 combinedDf.write.mode("overwrite").parquet("/path/to/save/result.parquet")
关于sc.textFile()的错误说明
sc.textFile("/data.sql").collect()无需传入参数,直接调用即可将RDD转为字符串数组,但该方法适用于分布式存储的大文件,本地小文件优先用scala.io.Source,避免不必要的分布式开销。
内容的提问来源于stack exchange,提问作者xBartegu
相关产品推荐
相关产品推荐

