从CSV构建Schema:遍历DataFrame生成StructField列表报错求助
问题分析与解决方案
首先咱们来拆解你遇到的问题:你写的代码里for循环没有使用yield,导致循环执行后返回的是Unit(相当于Java里的void),而你试图把这个Unit放到List()里,自然会触发「found: Unit required: StructField」的类型不匹配错误。另外代码里还有几个小细节需要修正,咱们一步步来:
先修正代码里的小错误
- 样例类里的
precision类型写错了,Scala里的整数类型是大写的Int,不是小写的int:
case class metadata_class(colname: String, datatype: String, length: Option[Int], precision: Option[Int])
- 读取CSV时的schema参数和类型转换写错了,修正后:
val foo = spark.read.format("csv") .option("delimiter", ",") .option("header", "true") .schema(Encoders.product[metadata_class].schema) .load("/path/to/file") .as[metadata_class] .toDF()
核心问题:正确构建StructField列表
你原来的代码试图用for循环直接生成List,但Scala中不带yield的for循环是执行副作用、返回Unit的,带yield才会收集循环的结果。另外更推荐用Scala集合的map操作(或Spark原生操作)来处理:
解法1:用Scala集合的map(适合小量元数据)
先把DataFrame转成Scala本地集合,再生成StructField列表:
import org.apache.spark.sql.types._ // 实现字符串类型转Spark DataType的函数 def getType(dataTypeStr: String, precision: Option[Int], length: Option[Int]): DataType = dataTypeStr match { case "string" => StringType case "int" => IntegerType case "decimal" => DecimalType(precision.getOrElse(10), length.getOrElse(2)) case "long" => LongType // 按需补充其他数据类型 } // 正确生成StructField列表 val sList: List[StructField] = foo.as[metadata_class].collect().map { m => StructField( name = m.colname, dataType = getType(m.datatype, m.precision, m.length), nullable = true // 可根据需求调整是否允许为空 ) }.toList
解法2:Spark原生操作(适合大量元数据,避免collect到本地)
如果元数据CSV数据量较大,不建议collect到本地,直接用DataFrame操作生成Schema:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 定义UDF将字符串类型转为Spark DataType val typeUdf = udf((dt: String, prec: Int, len: Int) => dt match { case "string" => StringType case "int" => IntegerType case "decimal" => DecimalType(prec, len) // 按需补充其他类型 }) // 生成最终的StructType val schema = foo.select( struct( col("colname"), typeUdf(col("datatype"), col("precision").cast("int"), col("length").cast("int")).alias("datatype") ).alias("field") ).agg(collect_list("field").alias("fields")) .selectExpr("transform(fields, x -> struct(x.colname as name, x.datatype as type, true as nullable)) as schema") .as[StructType] .first()
为什么原来的代码出错?
你写的代码:
val sList: List[StructField] = List( for (m <- foo.as[metadata_class].collect) { StructField(m.colname, getType(m.datatype)) })
这里的for循环没有加yield,执行后仅返回Unit,把它放到List()里得到的是List[Unit],和你声明的List[StructField]类型完全不匹配。正确的for循环写法应该是:
// 带yield的for循环写法 val sList: List[StructField] = for (m <- foo.as[metadata_class].collect) yield { StructField(m.colname, getType(m.datatype)) }
带yield的for循环会自动收集每次迭代生成的StructField,最终返回List[StructField],就不会有类型错误了。
内容的提问来源于stack exchange,提问作者Andrew
相关产品推荐
相关产品推荐

