Spark 3.2.1泛型场景下使用ExpressionEncoder报错问题
Spark Dataset泛型与编码器报错:Only expression encoders are supported for now
问题详情
在Spark Dataset中使用泛型与编码器时,触发错误:
org.apache.spark.SparkRuntimeException: Only expression encoders are supported for now
当前环境:Spark 3.2.1、Scala 2.12,将代码中的泛型替换为具体类型即可正常运行。
复现代码
import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.databind.node.ObjectNode import com.fasterxml.jackson.module.scala.DefaultScalaModule import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import org.apache.spark.sql.streaming.{StreamingQuery, Trigger} import org.apache.spark.sql.{Dataset, Encoder, SparkSession} abstract class SparkRunner[D]( deserializeFunc: Array[Byte] => Option[D]) (implicit spark: SparkSession, mapper: ObjectMapper) extends Serializable { import spark.implicits._ val linesOfBytes: Dataset[Array[Byte]] = spark.readStream .format("socket") .option("host", "localhost") .option("port", 9999) .load() .as[String].map(mapper.writeValueAsBytes) implicit val encD:Encoder[D] implicit val encOD:Encoder[Option[D]] val transformedDf = linesOfBytes .map(deserializeFunc) def start(): StreamingQuery = { transformedDf.writeStream //todo: Adopt config builder pattern for sink and src properties eventually. .format("console") .queryName("Query1") .trigger(Trigger.ProcessingTime(0)) .start() } } object SparkMain { def main(args: Array[String]): Unit = { implicit val spark: SparkSession = SparkSession .builder() .master("local") .appName("A") .getOrCreate() implicit val objMapper = new ObjectMapper() case class Payload(payload: ObjectNode) def deserialize(bytes: Array[Byte]): Option[String] = { objMapper.registerModule(DefaultScalaModule) val o = objMapper.readValue(bytes, classOf[Payload]) Some(objMapper.writeValueAsString(o.payload)) } try new SparkRunner( deserialize) { // override implicit val encD: Encoder[String] = org.apache.spark.sql.Encoders.kryo[String] // override implicit val encOD: Encoder[Option[String]] = org.apache.spark.sql.Encoders.kryo[Option[String]] override implicit val encD: Encoder[String] = ExpressionEncoder() override implicit val encOD: Encoder[Option[String]] = ExpressionEncoder() }.start().awaitTermination() finally { spark.close() } } }
//SparkMain.main(Array("")) // Uncomment if you're running this in spark-shell
解决方案
问题根源
泛型类SparkRunner[D]中,ExpressionEncoder()在编译时无法获取泛型D的具体类型信息,导致Spark无法生成有效的表达式编码器。换成具体类型时,编译器能直接推断类型,编码器正常工作。
方案1:用TypeTag保留泛型类型信息
修改SparkRunner类,引入Scala的TypeTag来保留泛型类型的元数据,让ExpressionEncoder能正确生成对应类型的编码器:
import scala.reflect.runtime.universe._ abstract class SparkRunner[D: TypeTag]( deserializeFunc: Array[Byte] => Option[D]) (implicit spark: SparkSession, mapper: ObjectMapper) extends Serializable { import spark.implicits._ val linesOfBytes: Dataset[Array[Byte]] = spark.readStream .format("socket") .option("host", "localhost") .option("port", 9999) .load() .as[String].map(mapper.writeValueAsBytes) // 通过TypeTag自动生成表达式编码器 implicit val encD: Encoder[D] = ExpressionEncoder[D]() implicit val encOD: Encoder[Option[D]] = ExpressionEncoder[Option[D]]() val transformedDf = linesOfBytes.map(deserializeFunc) def start(): StreamingQuery = { transformedDf.writeStream .format("console") .queryName("Query1") .trigger(Trigger.ProcessingTime(0)) .start() } }
修改后,子类实现时无需手动重写编码器,Scala会通过TypeTag保留泛型D的具体类型信息,确保ExpressionEncoder正常工作。
方案2:使用Kryo编码器替代
如果不需要表达式编码器的性能优化,也可以直接使用Kryo编码器,它对泛型的兼容性更好,不需要额外的类型信息。取消注释代码中的Kryo编码器部分即可:
override implicit val encD: Encoder[String] = org.apache.spark.sql.Encoders.kryo[String] override implicit val encOD: Encoder[Option[String]] = org.apache.spark.sql.Encoders.kryo[Option[String]]
注意:Kryo编码器的序列化效率略低于表达式编码器,但在泛型场景下更稳定。
内容的提问来源于stack exchange,提问作者Hemanth Gowda
相关产品推荐
相关产品推荐

