Scala Spark嵌套列表转DataFrame列数不匹配错误如何修复
问题原因
你定义的body是List[List[String]]类型,直接调用toDF()时,Spark会将每个子列表整体作为一行的单个值存储,最终生成的DataFrame仅包含1列,默认列名为value,当你尝试为这个只有1列的DataFrame指定7个列名时,就会触发列数不匹配的报错。
修复方案
有三种常用的实现方式,可根据场景选择:
方案1:转元组实现(简单快捷,列数固定时推荐)
将每行拆分得到的数组转为对应长度的元组,Spark可自动识别元组的每个元素为单独的列:
val spark = SparkSession.builder.appName("er").master("local").getOrCreate() import spark.implicits._ val erResponse = response.body.toString.split("\n") // 拆分表头并去除字段名前后空格 val columns = erResponse(0).split(",").map(_.trim) val body = erResponse.drop(1) .map(_.split(",").map(_.trim)) // 将7元素数组转为7元素元组 .map(arr => (arr(0), arr(1), arr(2), arr(3), arr(4), arr(5), arr(6))) .toList // 转DF时传入表头列名 val erDf = body.toDF(columns: _*) erDf.show()
方案2:构造Schema+Row实现(灵活适配列数变化)
如果后续列数可能调整,不想手动写固定长度的元组,可以通过自定义Schema+Row的方式创建DataFrame:
import org.apache.spark.sql.types.{StringType, StructField, StructType} import org.apache.spark.sql.Row val spark = SparkSession.builder.appName("er").master("local").getOrCreate() val erResponse = response.body.toString.split("\n") val columns = erResponse(0).split(",").map(_.trim) // 构造对应7个字符串类型字段的Schema val schema = StructType(columns.map(colName => StructField(colName, StringType, nullable = true))) // 把每行数据转成Row并构造RDD val rowRDD = spark.sparkContext.parallelize( erResponse.drop(1) .map(_.split(",").map(_.trim)) .map(arr => Row.fromSeq(arr)) ) val erDf = spark.createDataFrame(rowRDD, schema) erDf.show()
方案3:内置CSV读取器实现(无需手动拆分,自动处理格式)
因为你的原始数据本身是标准CSV格式,也可以直接把字符串转为Dataset后调用Spark内置的CSV读取器,无需自行处理拆分和表头逻辑:
val spark = SparkSession.builder.appName("er").master("local").getOrCreate() import spark.implicits._ // 将整个CSV内容转为字符串Dataset val csvDS = response.body.toString.split("\n").toSeq.toDS() // 直接用CSV读取器解析 val erDf = spark.read .option("header", "true") .option("trimValues", "true") .csv(csvDS) erDf.show()
内容的提问来源于stack exchange,提问作者megadarkfriend
相关产品推荐
相关产品推荐

