【问题标题】:How can I use case objects as a field in my case class and convert to Spark DataSet?如何将案例对象用作案例类中的字段并转换为 Spark DataSet?
【发布时间】:2021-07-14 20:42:06
【问题描述】:

我正在学习 spark-sql 并尝试在创建的 DataSet 中应用过滤器。我已经定义了一个简单的 Employee 案例类,它有 3 个字段,name、salary 和 dpt。

case class Employee( name: String, salary: Double, age: Int, dpt: Dept)

最后一个字段 dpt 定义如下:


 sealed trait Dept extends { val name: String }

  case object Accountability extends Dept { override val name = "AC"}
  case object Sales extends Dept { override val name = "S"}
  case object Finance extends Dept { override val name = "F"}
  case object Marketing extends Dept { override val name = "M"}
  case object Communication extends Dept { override val name = "C"}
  case object Reception extends Dept { override val name = "R"}
  case object HumanResource extends Dept { override val name = "HR"}

我试过用kryo编码器解决,但是不行。

 object DeptEncoders {
    implicit def deptEncoder : org.apache.spark.sql.Encoder[Dept] = org.apache.spark.sql.Encoders.kryo[Dept]
  }

【问题讨论】:

  • 也许你也需要Employee 的编码器

标签: scala apache-spark apache-spark-sql dataset


【解决方案1】:

基于文档here

import org.apache.spark.sql.Encoder

...
// conf is your org.apache.spark.SparkConf used to create your Spark Context
conf.registerKryoClasses(Array(classOf[Dept], classOf[Employee]))

...
implicit val encoder1:Encoder[Dept] = org.apache.spark.sql.Encoders.kryo[Dept]
implicit val encoder2:Encoder[Employee] = org.apache.spark.sql.Encoders.kryo[Employee]

...
val df = Seq(e1, e2).toDF()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-05-03
    • 2013-01-23
    • 1970-01-01
    • 2014-11-28
    • 2017-08-17
    • 2019-10-30
    • 1970-01-01
    • 2020-07-13
    相关资源
    最近更新 更多