【问题标题】:Spark No Encoder found for java.io.Serializable in Map[String, java.io.Serializable]Spark No Encoder found for java.io.Serializable in Map[String, java.io.Serializable]
【发布时间】:2018-11-15 23:12:03
【问题描述】:

我正在编写一个数据集非常灵活的 Spark 作业,它被定义为 Dataset[Map[String, java.io.Serializable]]

现在问题开始出现,spark 运行时抱怨No Encoder found for java.io.Serializable。我试过 kyro serde,仍然显示相同的错误消息。

我必须使用这种奇怪的数据集类型的原因是因为我每行都有灵活的字段。地图看起来像:

Map(
  "a" -> 1,
  "b" -> "bbb",
  "c" -> 0.1,
  ...
)

Spark 中是否有处理这种灵活数据集类型的方法?

编辑: 这是任何人都可以尝试的可靠代码。

import org.apache.spark.sql.{Dataset, SparkSession}

object SerdeTest extends App {
  val sparkSession: SparkSession = SparkSession
    .builder()
    .master("local[2]")
    .getOrCreate()


  import sparkSession.implicits._
  val ret: Dataset[Record] = sparkSession.sparkContext.parallelize(0 to 10)
    .map(
      t => {
        val row = (0 to t).map(
          i => i -> i.asInstanceOf[Integer]
        ).toMap

        Record(map = row)
      }
    ).toDS()

  val repartitioned = ret.repartition(10)


  repartitioned.collect.foreach(println)
}

case class Record (
                  map: Map[Int, java.io.Serializable]
                  )

上面的代码会给你错误编码器未找到:

Exception in thread "main" java.lang.UnsupportedOperationException: No Encoder found for java.io.Serializable
- map value class: "java.io.Serializable"
- field (class: "scala.collection.immutable.Map", name: "map")

【问题讨论】:

  • @thebluephantom 添加代码,可以直接在intelliJ中运行。

标签: apache-spark


【解决方案1】:

找到了答案,解决这个问题的一种方法是使用 Kyro serde 框架,代码更改非常少,只需要使用 Kyro 制作一个隐式编码器,并在需要序列化时将其带入上下文。

这是我工作的代码示例(可以直接在 IntelliJ 或等效 IDE 中运行):

import org.apache.spark.sql._

object SerdeTest extends App {
  val sparkSession: SparkSession = SparkSession
    .builder()
    .master("local[2]")
    .getOrCreate()


  import sparkSession.implicits._

  // here is the place you define your Encoder for your custom object type, like in this case Map[Int, java.io.Serializable]
  implicit val myObjEncoder: Encoder[Record] = org.apache.spark.sql.Encoders.kryo[Record]
  val ret: Dataset[Record] = sparkSession.sparkContext.parallelize(0 to 10)
    .map(
      t => {
        val row = (0 to t).map(
          i => i -> i.asInstanceOf[Integer]
        ).toMap

        Record(map = row)
      }
    ).toDS()

  val repartitioned = ret.repartition(10)


  repartitioned.collect.foreach(
    row => println(row.map)
  )
}

case class Record (
                  map: Map[Int, java.io.Serializable]
                  )

这段代码会产生预期的结果:

Map(0 -> 0, 5 -> 5, 1 -> 1, 2 -> 2, 3 -> 3, 4 -> 4)
Map(0 -> 0, 1 -> 1, 2 -> 2)
Map(0 -> 0, 5 -> 5, 1 -> 1, 6 -> 6, 2 -> 2, 7 -> 7, 3 -> 3, 4 -> 4)
Map(0 -> 0, 1 -> 1)
Map(0 -> 0, 1 -> 1, 2 -> 2, 3 -> 3, 4 -> 4)
Map(0 -> 0, 1 -> 1, 2 -> 2, 3 -> 3)
Map(0 -> 0)
Map(0 -> 0, 5 -> 5, 1 -> 1, 6 -> 6, 2 -> 2, 3 -> 3, 4 -> 4)
Map(0 -> 0, 5 -> 5, 10 -> 10, 1 -> 1, 6 -> 6, 9 -> 9, 2 -> 2, 7 -> 7, 3 -> 3, 8 -> 8, 4 -> 4)
Map(0 -> 0, 5 -> 5, 1 -> 1, 6 -> 6, 9 -> 9, 2 -> 2, 7 -> 7, 3 -> 3, 8 -> 8, 4 -> 4)
Map(0 -> 0, 5 -> 5, 1 -> 1, 6 -> 6, 2 -> 2, 7 -> 7, 3 -> 3, 8 -> 8, 4 -> 4)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-07-28
    • 1970-01-01
    • 2020-06-28
    • 1970-01-01
    • 2015-01-26
    相关资源
    最近更新 更多