Spark读取Parquet中存储的Array[Byte]数据报类型转换错误如何解决
问题产生原因
Spark SQL 内部处理ArrayType(ByteType)类型字段时,默认不会返回 Scala/Java 原生的字节数组类型Array[Byte](即报错信息中的[B类型),而是使用scala.collection.mutable.WrappedArray对数组元素做封装:
- 写入Parquet时,你传入的
Array[Byte]会被Spark自动封装为WrappedArray落盘 - 读取时直接调用
row.getAs[Array[Byte]]相当于将没有继承关系的WrappedArray强转为原生字节数组,自然触发类型转换异常
解决方案
可根据你的使用场景任选以下任意一种方案:
- 方案1:手动转换
WrappedArray为原生数组
先读取到Spark返回的WrappedArray实例,再调用toArray方法转成原生字节数组,适配所有Spark版本:import scala.collection.mutable.WrappedArray val wrappedArr = row.getAs[WrappedArray[Byte]]("byteArrayObject") val byteArr = wrappedArr.toArray - 方案2:使用Row原生提供的字节数组读取方法
Spark的Row接口自带getByteArray专用方法,内部已经完成类型适配,不需要手动处理封装类型:// 先获取字段的索引位置,再调用getByteArray val byteArr = row.getByteArray(row.fieldIndex("byteArrayObject")) - 方案3:通过Dataset强类型映射自动转换
提前定义对应的数据结构case class,将读取到的DataFrame转为强类型Dataset,Spark内置的编码器会自动完成WrappedArray到原生字节数组的转换:// 定义和Parquet字段完全匹配的case class case class ParquetRecord(byteArrayObject: Array[Byte], /* 补充其他字段定义 */) // 读取Parquet后直接转为Dataset val ds = spark.read.parquet("你的parquet路径").as[ParquetRecord] // 后续直接访问ds的byteArrayObject属性即为原生Array[Byte]类型
内容的提问来源于stack exchange,提问作者kailing
相关产品推荐
相关产品推荐

