使用Apache ORC 1.8(Kotlin)写入字符串到ORC文件失败求助
解决Apache ORC 1.8写入字符串后Spark读取为空的问题
问题诊断
空值问题的核心是错误使用BytesColumnVector的赋值逻辑:ORC的字节列向量并非直接通过vector[row]赋值字节数组就能被正确识别,它依赖start(数据起始偏移)、length(有效数据长度)和isNull(是否为空)三个辅助数组标记每行数据状态。直接赋值vector[row]时,这些辅助字段未被初始化,导致读取时无法解析有效内容,最终显示为空。
修复方案
使用BytesColumnVector提供的setVal方法自动处理字节存储与辅助字段设置,或者手动初始化start、length和isNull字段。
修复后的完整代码
import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.Path import org.apache.hadoop.hive.ql.exec.vector.BytesColumnVector import org.apache.orc.OrcFile import org.apache.orc.TypeDescription import org.apache.orc.Writer fun main() { val conf = Configuration() val schema = TypeDescription.fromString("struct<x:string,y:string>") val writer = OrcFile.createWriter( Path("./my-file.orc"), OrcFile.writerOptions(conf) .setSchema(schema) ) val batch = schema.createRowBatch() val x = batch.cols[0] as BytesColumnVector val y = batch.cols[1] as BytesColumnVector for (r in 0..99) { val row = batch.size++ val bytes = r.toString().toByteArray(Charsets.UTF_8) // 用setVal自动处理辅助字段 x.setVal(row, bytes) y.setVal(row, bytes) if (batch.size == batch.maxSize) { writer.addRowBatch(batch) batch.reset() } } if (batch.size != 0) { writer.addRowBatch(batch) batch.reset() } writer.close() }
手动初始化辅助字段的替代写法
如果不想使用setVal,可手动设置相关字段:
val bytes = r.toString().toByteArray(Charsets.UTF_8) x.vector[row] = bytes x.start[row] = 0 x.length[row] = bytes.size x.isNull[row] = false y.vector[row] = bytes y.start[row] = 0 y.length[row] = bytes.size y.isNull[row] = false
验证结果
重新运行写入代码后,用Spark读取即可正常显示数据:
val df = spark.read.format("orc").load("./my-file.orc") df.show() df.printSchema()
输出示例:
+---+---+ | x| y| +---+---+ | 0| 0| | 1| 1| | 2| 2| | ...| ...| +---+---+
内容的提问来源于stack exchange,提问作者jake wong
相关产品推荐
相关产品推荐

