【问题标题】:How to convert bytes from Kafka to their original object?如何将字节从 Kafka 转换为其原始对象?
【发布时间】:2017-11-01 03:50:13
【问题描述】:

我正在从 Kafka 获取数据,然后使用默认解码器反序列化 Array[Byte],之后我的 RDD 元素看起来像 (null,[B@406fa9b2)、(null,[B@21a9fe0) 但我想要我的原始数据有一个架构,那么我该怎么做实现这个?

我以 Avro 格式序列化消息。

【问题讨论】:

    标签: apache-spark apache-kafka spark-streaming spark-avro


    【解决方案1】:

    您必须使用适当的反序列化器解码字节,例如字符串或您的自定义对象。

    如果您不进行解码,您会得到 [B@406fa9b2,这只是 Java 中字节数组的文本表示。

    Kafka 对消息的内容一无所知,因此它将字节数组从生产者传递给消费者。

    在 Spark Streaming 中,您必须对键和值使用序列化程序(引用 KafkaWordCount example):

    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      "org.apache.kafka.common.serialization.StringSerializer")
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
      "org.apache.kafka.common.serialization.StringSerializer")
    

    使用上述序列化程序,您将获得DStream[String],因此您可以使用RDD[String]。

    但是,如果您想直接将字节数组反序列化为自定义类,则必须编写自定义 Serializer(这是 Kafka 特有的,与 Spark 无关)。

    我建议使用带有固定架构的 JSON 或 Avro(使用Kafka, Spark and Avro - Part 3, Producing and consuming Avro messages 中描述的解决方案)。


    在Structured Streaming 中,管道可能如下所示:

    val fromKafka = spark.
      readStream.
      format("kafka").
      option("subscribe", "topic1").
      option("kafka.bootstrap.servers", "localhost:9092").
      load.
      select('value cast "string") // <-- conversion here
    

    【讨论】:

    • 那么,如何在 Spark Structured Streaming 中将带/不带模式注册表的 Avro Kafka 消息转换为原始对象?
    • 你必须知道原来的对象,例如使用map操作符。还没有from_avro(如果有的话),就像我们对带有from_json 的JSON 所做的那样。
    • 我使用 KafkaAvroDeserializer 将 Array[Byte] 映射到我的 Avro 对象,但它说“无法找到存储在数据集中的类型的编码器”。然后我提供编码器作为隐式 def toEncoded(o: Zhima): Array[Byte] = o.toByteBuffer.array() implicit def fromEncoded(e: Array[Byte]): Zhima = valueDeserializer.deserialize(kafkaConsumeTopicName, e).asInstanceOf [Zhima] 但是也出现了同样的错误,那怎么解决呢?
    • 然后我使用自定义 UDF 来解析 avro 消息,现在它报告“线程“主”java.lang.UnsupportedOperationException 中的异常:不支持 org.apache.avro.generic.GenericRecord 类型的架构” .我是否需要将 Avro 对象转换为 scala 案例类或 java pojo? spark.udf.register("deserialize", (topic: String, bytes: Array[Byte]) => MyDeserializerWrapper.deser.deserialize(topic, bytes).asInstanceOf[GenericRecord] )
    猜你喜欢
    • 2021-03-12
    • 2014-07-24
    • 1970-01-01
    • 2016-04-11
    • 2015-12-19
    • 1970-01-01
    • 2011-05-01
    • 1970-01-01
    相关资源
    最近更新 更多