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