如何为带上界约束的泛型Product类配置Spark Encoder
如何为带上界约束的泛型
Product类使用对应的Product编码器 以下代码无法完成编译:
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本身继承自Product,理论上可作为Dataset的类型使用,其定义如下:
abstract class KafkaValueContract(val metadata: Metadata, val changes: Changes) extends Product with Serializable
泛型参数VALUE为任意继承KafkaValueContract的样例类,示例实现如下:
case class PlaceholderDataContract(override val metadata: PlaceholderMetadata, override val changes: PlaceholderChanges) extends KafkaValueContract(metadata, changes)
实际编译时报错提示不存在KafkaMessage[KafkaKeyContract, VALUE]对应的编码器。预期逻辑为:VALUE是继承自KafkaValueContract(该类已继承Product)的任意样例类,应当可自动生成对应Encoder,实际报错信息如下:
[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] ^
问题更新
若在类签名中为泛型参数添加TypeTag上下文绑定,显式告知Scala该泛型对应具体类型,编译器可识别到隐式Product编码器:
class EnrichmentParser[VALUE <: KafkaValueContract : TypeTag](typeParser: TypeParser[VALUE]) extends Serializable {
但代码运行时会抛出如下反射异常:
type _$1 is not a class scala.ScalaReflectionException: type _$1 is not a class
解决方案
Spark的隐式Encoder生成逻辑针对具体静态类型在编译时生成,仅通过泛型上界声明VALUE <: Product、绑定TypeTag无法满足要求:泛型擦除后,运行时反射只能拿到存在类型标记,无法获取具体样例子类的字段结构,才会抛出上述反射异常。
正确处理方式是直接在泛型参数上增加Encoder上下文绑定,要求实例化EnrichmentParser时(此时VALUE为确定的具体样例类),由调用处的Spark隐式推导生成对应类型的Encoder并传入,内部再组合生成KafkaMessage的通用编码器即可:
import org.apache.spark.sql.Encoder import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder // 给VALUE增加Encoder上下文绑定,替换原有的TypeTag绑定 class EnrichmentParser[VALUE <: KafkaValueContract : Encoder](typeParser: TypeParser[VALUE]) extends Serializable { // 基于已有的KafkaKeyContract、VALUE的Encoder,直接生成KafkaMessage对应的Encoder implicit private val kafkaMsgEncoder: Encoder[KafkaMessage[KafkaKeyContract, VALUE]] = ExpressionEncoder[KafkaMessage[KafkaKeyContract, VALUE]] 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() } }
如果KafkaKeyContract本身也是抽象类、对应不同子类实现,只需按照相同逻辑给该类型也增加对应的Encoder上下文绑定即可。
内容的提问来源于stack exchange,提问作者YFl
相关产品推荐
相关产品推荐

