Spark Scala中将DataFrame转换为DataSet遇到的类型匹配问题
嘿,这个问题我太熟了!你之前把sparse_rep以Map[Int,Double]类型写入MongoDB,但Spark默认会把Map转成Struct存储,所以读回来的时候这个字段变成了Struct类型,和你样例类里的Map类型不匹配,才会抛出那个错误:
Caused by: org.apache.spark.sql.AnalysisException: need a map field but got struct<1:double,10:double,11:double,12:double,13:double,14:double,15:double,17:double,18:double,2:double,20:double,21:double,22:double,23:double,24:double,25:double,26:double,27:double,3:double,4:double,7:double,9:double>;
别担心,咱们一步步把这个Struct转回Map就行:
步骤1:获取Struct的所有字段名
首先得拿到sparse_rep这个Struct里的所有字段名,这些就是原来Map的key(只是现在是字符串形式的数字):
val structFields = db.select("sparse_rep.*").columns
步骤2:把Struct转换成Map
用Spark的map_from_arrays函数,把Struct的字段名转成Int类型的key,字段值作为value,组合成Map:
import org.apache.spark.sql.functions._ // 把字段名转成Int类型的数组作为Map的key val keysArray = array(structFields.map(field => lit(field.toInt)): _*) // 提取Struct中每个字段的值作为Map的value val valuesArray = array(structFields.map(field => col(s"sparse_rep.${field}")): _*) // 替换原来的sparse_rep字段为Map类型 val dfWithMap = db.withColumn("sparse_rep", map_from_arrays(keysArray, valuesArray))
步骤3:转换为样例类
现在sparse_rep已经是Map[Int,Double]类型了,就可以正常转换成你定义的blogRow样例类了:
case class blogRow(_id:String, id:Int, sparse_rep:Map[Int,Double],title:String) val blogRowEncoder = Encoders.product[blogRow] val blogDS = dfWithMap.as[blogRow](blogRowEncoder)
额外提示:避免后续再踩坑
如果之后还要把数据写回MongoDB,记得加个配置,阻止Spark把Map转成Struct,这样下次读的时候就直接是Map类型了:
import com.mongodb.spark.config.WriteConfig val writeConfig = WriteConfig(Map("spark.mongodb.output.convertMapToStruct" -> "false")) blogDS.write.mode("overwrite").format("mongo").options(writeConfig.asOptions).save()
内容的提问来源于stack exchange,提问作者Sajeed

