【问题标题】:Decoding Java enums/custom non case classes using Structured Spark Streaming使用结构化 Spark 流解码 Java 枚举/自定义非案例类
【发布时间】:2018-10-27 20:35:51
【问题描述】:

我正在尝试使用 Spark 2.1.1 中的结构化流来读取 Kafka 并解码 Avro 编码的消息。我有一个根据 this question.

定义的 UDF
val sr = new CachedSchemaRegistryClient(conf.kafkaSchemaRegistryUrl, 100)
val deser = new KafkaAvroDeserializer(sr)

val decodeMessage = udf { bytes:Array[Byte] => deser.deserialize("topic.name", bytes).asInstanceOf[DeviceRead] }

val topic = conf.inputTopic
val df = session
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", conf.kafkaServers)
    .option("subscribe", topic)
    .load()

df.printSchema()

val result = df.selectExpr("CAST(key AS STRING)", """decodeMessage($"value") as "value_des"""")

val query = result.writeStream
    .format("console")
    .outputMode(OutputMode.Append())
    .start()

但是我得到以下失败。

Exception in thread "main" java.lang.UnsupportedOperationException: Schema for type DeviceRelayStateEnum is not supported

在这一行失败

val decodeMessage = udf { bytes:Array[Byte] => deser.deserialize("topic.name", bytes).asInstanceOf[DeviceRead] }

另一种方法是为我拥有的自定义类定义编码器

implicit val enumEncoder = Encoders.javaSerialization[DeviceRelayStateEnum]
implicit val messageEncoder = Encoders.product[DeviceRead]

但是在注册 messageEncoder 时失败并出现以下错误。

Exception in thread "main" java.lang.UnsupportedOperationException: No Encoder found for DeviceRelayStateEnum
- option value class: "DeviceRelayStateEnum"
- field (class: "scala.Option", name: "deviceRelayState")
- root class: "DeviceRead"
    at org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:602)
    at org.apache.spark.sql.catalyst.ScalaReflection$.org$apache$spark$sql$catalyst$ScalaReflection$$serializerFor(ScalaReflection.scala:476)
    at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:596)
    at org.apache.spark.sql.catalyst.ScalaReflection$$anonfun$9.apply(ScalaReflection.scala:587)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.immutable.List.foreach(List.scala:381)
    at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)

当我尝试在 load() 之后使用 map 执行此操作时,我收到以下编译错误。

val result = df.map((bytes: Row) => deser.deserialize("topic", bytes.getAs[Array[Byte]]("value")).asInstanceOf[DeviceRead])

Error:(76, 26) not enough arguments for method map: (implicit evidence$6: org.apache.spark.sql.Encoder[DeviceRead])org.apache.spark.sql.Dataset[DeviceRead].
Unspecified value parameter evidence$6.
      val result = df.map((bytes: Row) => deser.deserialize("topic", bytes.getAs[Array[Byte]]("value")).asInstanceOf[DeviceRead])
Error:(76, 26) Unable to find encoder for type stored 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.
      val result = df.map((bytes: Row) => deser.deserialize("topic", bytes.getAs[Array[Byte]]("value")).asInstanceOf[DeviceRead])

这是否意味着我不能对 Java 枚举使用结构化流?而且它只能与原语或案例类一起使用?

我阅读了一些相关的问题123 围绕这个问题,似乎有可能为一个类指定一个自定义编码器,即 UDT 在 2.1 中被删除并且没有添加新功能。

我们将不胜感激。

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-structured-streaming


    【解决方案1】:

    认为一般而言,您可能对当前版本的结构化流(和 Spark SQL)要求过高。

    我还不能完全理解如何以所谓的更专业的方式处理丢失编码器的问题,但是当您尝试创建枚举的Dataset 时会遇到同样的问题。这可能还没有被简单地支持。

    Structured Streaming 只是 Spark SQL 之上的一个流库,并将其用于序列化-反序列化 (SerDe)。

    为了使故事简短并让您继续前进(直到您找到更好的方法),我建议避免在您用来表示数据集架构的业务对象中使用枚举。

    所以,我建议按照以下方式做一些事情:

    val decodeMessage = udf { bytes:Array[Byte] =>
      val dr = deser.deserialize("topic.name", bytes).asInstanceOf[DeviceRead]
    
      // do additional transformation here so you use a custom streaming-specific class
      // Here I'm using a simple tuple to hold what might be relevant
      // You could create a case class instead to have proper names
      (dr.id, dr.value)
    }
    

    【讨论】:

    • 问题是我们在 Spark 1.6 中使用了枚举和上面定义的格式,使用了直接流式处理方法。所以我们不能在不影响其他应用程序的情况下使用其他东西。如果我能找到其他东西,我会四处寻找。
    • 临时映射仅用于使用 Spark(SQL 和结构化流)处理数据,外部应用程序甚至无法知道内部“打嗝”。直接流方法不同,因为它不使用 Spark SQL 的编码器。我很乐意听到更多关于您的担忧的信息,以便提供更好的解决方案。介意详细说明更新您的问题吗?
    猜你喜欢
    • 1970-01-01
    • 2013-09-04
    • 1970-01-01
    • 1970-01-01
    • 2022-01-18
    • 2012-10-04
    • 1970-01-01
    • 2015-08-12
    • 1970-01-01
    相关资源
    最近更新 更多