Spark转换为RDD操作失败,无法基于CSV读取的表名字段创建DataFrame
Spark 代码错误排查与实现方案
原代码核心错误点
- 调用
spark.createDataFrame()时直接传入字符串类型的tab1不符合方法参数要求,该方法需要传入RDD、Seq等可序列化的集合类数据,不能直接接收单个字符串 - 重复调用
row.mkString(",").split(",")完全冗余,CSV读取得到的Row对象可以直接通过索引获取对应列值,强行转字符串拆分还可能因为列值本身包含逗号导致取值错误 - 代码中
tab2和tab3都取了索引为1的列,属于笔误,tab3应该取索引为2的列 - 没有将执行SQL得到的结果追加表名列,也没有对多个结果做合并操作
正确实现代码
import org.apache.spark.sql.functions.lit // 读取CSV文件,如果你的CSV有表头可以添加.option("header", "true")配置 val df1 = spark.read.format("csv").load("c:\\file.csv") // 收集配置行数据,若配置行数极多不建议使用collect,避免Driver内存溢出 val configRows = df1.collect() // 遍历所有配置行,生成结果DataFrame后合并 val resultDF = configRows.map(row => { // 直接从Row中取对应列值,无需转字符串拆分 val tabName = row.getAs[String](0) val sqlTab2 = row.getAs[String](1) val sqlTab3 = row.getAs[String](2) // 执行两个查询语句,分别追加固定值表名列Col0 val dfTab2 = spark.sql(sqlTab2).withColumn("Col0", lit(tabName)) val dfTab3 = spark.sql(sqlTab3).withColumn("Col0", lit(tabName)) // 合并同一个配置行的两个查询结果 dfTab2.unionByName(dfTab3) }).reduce(_ unionByName _) // 调整列顺序为要求的Col0、Col1、Col2,可根据实际业务调整 val finalDF = resultDF.select("Col0", "Col1", "col2") // 可直接创建临时视图供后续Spark SQL查询使用 finalDF.createOrReplaceTempView("result_table")
注意事项
- 执行SQL前请确保
sqlTab2和sqlTab3的输出列结构一致,否则unionByName会报错 - 如果CSV存在特殊分隔符、引号转义等规则,读取CSV时可以补充对应的option配置适配
- 若配置行数特别多,避免使用collect把所有数据拉到Driver端,可以改用广播变量等分布式方案实现
内容的提问来源于stack exchange,提问作者VnS
相关产品推荐
相关产品推荐

