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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 00:17:36