Spark Encoder for generic T with upper bound

Viewed 22

How to work with a product encoder for a generic upper-bounded Product class?

The following code doesn't compile:

class EnrichmentParser[VALUE <: KafkaValueContract](typeParser: TypeParser[VALUE]) extends Serializable {

  private def parseKey(row: Row): KafkaKeyContract = ...

  private def parseValue(row: Row): VALUE = typeParser.parse(row)

  private def parseRow(row: Row): KafkaMessage[KafkaKeyContract, VALUE] = {
    val key = parseKey(row)
    val value = parseValue(row)
    KafkaMessage(Some(key), value)
  }

  def parse(df: DataFrame)(implicit spark: SparkSession): DataFrame = {
    import spark.implicits._

    df.map(row => parseRow(row)).toDF()

  }

KafkaValueContract naturally extends Product so we can use it as a Dataset type:

abstract class KafkaValueContract(val metadata: Metadata,
                                  val changes: Changes) extends Product with Serializable

VALUE is whatever case class extending KafkaValueContract, for example:

case class PlaceholderDataContract(override val metadata: PlaceholderMetadata,
                                   override val changes: PlaceholderChanges) extends KafkaValueContract(metadata, changes)

However, the compiler complains that there is no encoder for KafkaMessage[KafkaKeyContract, VALUE], which I expected to have since VALUE is any other (case) class that extends KafkaValueContract which in turn extends Product:

[error] ... KafkaMessage[KafkaKeyContract,VALUE]. An implicit Encoder[KafkaMessage[KafkaKeyContract,VALUE]] is needed to store KafkaMessage[KafkaKeyContract,VALUE] instances in a Dataset. Primitive types (Int, String, etc) and Product types (case classes) are supported by importing spark.implicits._  Support for serializing other types will be added in future releases.
[error]     df.map(row => parseRow(row)).toDF()
[error]           ^

Thanks.


UPDATE:

If I add TypeTag to the class signature in order to explicitly tell Scala that the class is concrete it finds the implicit product encoder:

class EnrichmentParser[VALUE <: KafkaValueContract : TypeTag](typeParser: TypeParser[VALUE]) extends Serializable {

However, it throws a reflection exception:

type _$1 is not a class
scala.ScalaReflectionException: type _$1 is not a class
0 Answers
Related