保存含Vector类型的Spark DataFrame为ORC后读取Schema变更问题
这个问题我之前也碰到过,本质是Spark ML的Vector类型和ORC格式的兼容性问题——ORC并不原生支持Spark自定义的Vector类型,所以写入时会被序列化为ORC原生的struct结构,读取时就没法自动还原成Vector了。
先明确你遇到的现象:
原DataFrame的Schema是:
root |-- label: double (nullable = true) |-- features: vector (nullable = true)
而从ORC加载后的Schema会变成:
root |-- label: double (nullable = true) |-- features: struct (nullable = true) | |-- values: array (nullable = true) | | |-- element: double (containsNull = true)
下面给你两种可靠的解决方案:
方案一:加载后手动将Struct转换回Vector
我们可以用UDF把ORC解析出来的struct重新转换成Spark的Vector类型:
- 先导入需要的依赖:
import org.apache.spark.ml.linalg.Vectors import org.apache.spark.sql.functions.{col, udf}
- 定义转换UDF:
val structToVector = udf { row: org.apache.spark.sql.Row => // 从struct中取出values数组,转成Vector Vectors.dense(row.getAs[Seq[Double]]("values").toArray) }
- 应用UDF修复Schema:
val fixedDF = newDF.withColumn("features", structToVector(col("features"))) fixedDF.printSchema
执行后,fixedDF的Schema就和原DataFrame完全一致了。
方案二:写入前先将Vector转成数组,加载后再还原
如果想避免加载时的转换操作,可以在写入ORC前先把Vector转成ORC原生支持的数组类型,加载后再转回Vector:
- 写入前处理DataFrame:
// 将Vector列转成array<double>类型 val dfForOrc = df.withColumn("features", col("features").cast("array<double>")) dfForOrc.write.mode(SaveMode.Overwrite).orc("/some/path")
- 加载后还原成Vector:
val newDF = spark.read.orc("/some/path") val arrayToVector = udf { arr: Array[Double] => Vectors.dense(arr) } val fixedDF = newDF.withColumn("features", arrayToVector(col("features")))
额外说明
为什么会出现这个问题?因为Spark ML的Vector是Spark自定义的复杂类型,ORC格式的元数据里没有对应的类型标识,所以Spark只能将其拆解成ORC能识别的struct结构存储。读取时,Spark默认会按照ORC的struct结构解析,而不会自动关联到Spark的Vector类型,所以就出现了Schema不一致的情况。
如果你的场景允许,也可以考虑使用Parquet格式——Parquet对Spark自定义类型的支持更好,保存和加载Vector列时Schema不会发生变化,这可能是更省心的选择。
内容的提问来源于stack exchange,提问作者Neelesh Sambhajiche
相关产品推荐
相关产品推荐

