Encode an ADT / sealed trait hierarchy into Spark DataSet column
Asked Answered
C

3

14

If I want to store an Algebraic Data Type (ADT) (ie a Scala sealed trait hierarchy) within a Spark DataSet column, what is the best encoding strategy?

For example, if I have an ADT where the leaf types store different kinds of data:

sealed trait Occupation
case object SoftwareEngineer extends Occupation
case class Wizard(level: Int) extends Occupation
case class Other(description: String) extends Occupation

Whats the best way to construct a:

org.apache.spark.sql.DataSet[Occupation]
Christos answered 8/12, 2016 at 1:3 Comment(0)
C
11

TL;DR There is no good solution right now, and given Spark SQL / Dataset implementation, it is unlikely there will be one in the foreseeable future.

You can use generic kryo or java encoder

val occupation: Seq[Occupation] = Seq(SoftwareEngineer, Wizard(1), Other("foo"))
spark.createDataset(occupation)(org.apache.spark.sql.Encoders.kryo[Occupation])

but is hardly useful in practice.

UDT API provides another possible approach as for now (Spark 1.6, 2.0, 2.1-SNAPSHOT) it is private and requires quite a lot boilerplate code (you can check o.a.s.ml.linalg.VectorUDT to see example implementation).

Cattier answered 11/12, 2016 at 2:57 Comment(1)
why is kryo hardly useful in practice? Is it because we have to specify kryo serializer after every transformation?Grime
S
2

I have once dived deeply into the subject and created a repo showcasing all the approaches I have found could be useful.

Link: https://github.com/atais/spark-enum

Generally, zero323 is right, but you might find it useful to understand the full picture.

Stain answered 11/8, 2020 at 16:21 Comment(0)
A
-1

fyi for anyone attempting this: https://mcmap.net/q/831215/-scala-dataset-with-case-class-inheritance, frameless encoder derivation works.

The default Spark Scala encoding works on concrete product's only. frameless' injections and encoder derivation let you represent different encodings e.g. ADTs

Affer answered 18/4 at 10:3 Comment(1)
Please add all clarification to your answer instead of posting links to other questionsScorper

© 2022 - 2024 — McMap. All rights reserved.