Spark 3.5中RowEncoder替代及Row与InternalRow序列化方案问询
Spark 3.5中实现Row与InternalRow互转的序列化/反序列化器替代方案
旧版本(Spark 3.1)代码
Spark 3.1中用于创建序列化器、反序列化器的代码:
val encoder = Encoders.row(schema).resolveAndBind() val serializer = encoder.createSerializer() val deserializer = encoder.createDeserializer()
Spark 3.5适配方案
Spark 3.5移除了RowEncoder.apply(),且Encoders.row(schema)返回的Encoder[Row]不再直接暴露createSerializer()和createDeserializer()方法,可通过以下两种方式实现Row与InternalRow的互转:
方案1:基于ExpressionEncoder标准API实现(推荐)
Encoders.row(schema)本质是ExpressionEncoder[Row],可通过绑定Spark环境后获取序列化/反序列化逻辑:
import org.apache.spark.sql.{Encoders, SparkSession} import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.types.StructType import org.apache.spark.sql.Row // 假设已初始化spark: SparkSession和schema: StructType val encoder = Encoders.row(schema).asInstanceOf[ExpressionEncoder[Row]] // 将编码器绑定到当前Spark执行环境 val boundEncoder = encoder.resolveAndBind(spark.sessionState.conf) // Row转InternalRow的序列化函数 val rowToInternalRow: Row => InternalRow = boundEncoder.createSerializer() // InternalRow转Row的反序列化函数 val internalRowToRow: InternalRow => Row = boundEncoder.createDeserializer()
方案2:使用内置转换工具(内部API,需注意兼容性)
如果仅需简单互转,可借助Spark内部的RowConverter工具(注:内部API不保证跨版本兼容):
import org.apache.spark.sql.catalyst.util.RowConverter import org.apache.spark.sql.types.StructType import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.Row // 假设已有schema: StructType val converter = RowConverter.createConverter(schema) // Row转InternalRow val rowToInternalRow: Row => InternalRow = (row: Row) => converter(row).asInstanceOf[InternalRow] // InternalRow转Row val internalRowToRow: InternalRow => Row = (internalRow: InternalRow) => Row.fromSeq(internalRow.toSeq(schema))
注意事项
- 方案1基于Spark标准API,稳定性更高,优先推荐;
- 确保转换的Row/InternalRow结构与目标
schema完全匹配,否则会触发类型转换异常。
内容的提问来源于stack exchange,提问作者Joy Cheng
相关产品推荐
相关产品推荐

