You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 01:54:54