Scala Spark中如何将Vector[String]转换为DataFrame?
问题分析与解决方案
错误原因
你的代码中使用split("\s+")按空格分割每行数据,但原始数据是逗号分隔的CSV格式。这会导致每行分割后得到的数组长度不匹配schema定义的3个字段(比如表头行分割后仅1个元素),直接触发ArrayIndexOutOfBoundsException。
修正方案
方案1:修复原有代码逻辑
调整分割符为逗号,同时基于表头生成贴合业务的schema:
import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType} import org.apache.spark.sql.Row // 分离表头和数据行 val header = dataInVectorString.head.split(",") // 分割数据行,可选trim处理可能存在的空格 val dataRows = dataInVectorString.tail.map(_.split(",").map(_.trim)) // 转换为RDD[Row] val rowRdd = sparkSession.sparkContext.parallelize(dataRows).map(Row.fromSeq(_)) // 基于表头构建schema,指定对应数据类型 val schema = StructType( header.zipWithIndex.map { case (colName, idx) => val colType = idx match { case 0 => StringType // user_uid 为字符串类型 case 1 => IntegerType // score 为整数类型 case 2 => IntegerType // value 为整数类型 } StructField(colName, colType, nullable = true) } ) // 创建DataFrame val df = sparkSession.createDataFrame(rowRdd, schema) df.show()
方案2:使用Spark原生CSV读取API(更简洁)
利用Spark内置的CSV解析能力,无需手动处理分割和Row转换:
import sparkSession.implicits._ val df = sparkSession.read .option("header", "true") // 第一行作为表头 .option("inferSchema", "true") // 自动推断字段数据类型 .csv(dataInVectorString.toDS()) // 将Vector转为Dataset[String]后读取 df.show()
内容的提问来源于stack exchange,提问作者user16619164
相关产品推荐
相关产品推荐

