【问题标题】:How to deserialize Avro messages from Kafka in Flink (Scala)?如何在 Flink(Scala)中反序列化来自 Kafka 的 Avro 消息?
【发布时间】:2019-08-03 08:37:36
【问题描述】:

我正在将来自 Kafka 的消息读入 Flink Shell (Scala),如下:

scala> val stream = senv.addSource(new FlinkKafkaConsumer011[String]("topic", new SimpleStringSchema(), properties)).print()
warning: there was one deprecation warning; re-run with -deprecation for details
stream: org.apache.flink.streaming.api.datastream.DataStreamSink[String] = org.apache.flink.streaming.api.datastream.DataStreamSink@71de1091

在这里,我使用 SimpleStringSchema() 作为反序列化器,但实际上消息具有另一个 Avro 架构(例如 msg.avsc)。如何基于这种不同的 Avro 架构 (msg.avsc) 创建反序列化器,以反序列化传入的 Kafka 消息?

我无法在 Scala 中找到任何代码示例或教程,因此任何输入都会有所帮助。看来我可能需要扩展和实现

org.apache.flink.streaming.util.serialization.DeserializationSchema

用于解码消息,但我不知道该怎么做。任何教程或说明都会有很大帮助。因为,我不想进行任何自定义处理,而只是按照 Avro 架构 (msg.avsc) 解析消息,所以任何快速的方法都会非常有帮助。

【问题讨论】:

标签: scala apache-kafka deserialization apache-flink avro


【解决方案1】:

我在 java 中找到了 AvroDeserializationSchema 类的示例

https://github.com/okkam-it/flink-examples/blob/master/src/main/java/org/okkam/flink/avro/AvroDeserializationSchema.java

代码 sn-p:

如果您想反序列化为特定的案例类,请使用new FlinkKafkaConsumer011[case_class_name]new AvroDeserializationSchema[case_class_name](classOf[case_class_name]

val stream = env .addSource(new FlinkKafkaConsumer011[DeviceData]
 ("test", new AvroDeserializationSchema[case_class_name](classOf[case_class_name]), properties))

如果您使用 Confluent 的模式注册表,那么首选的解决方案是使用 Confluent 提供的 Avro serde。我们只需调用 deserialize() 并且要使用的最新版本的 Avro 模式的解析是在幕后自动完成的,并且不需要字节操作。

在 scala 中如下所示。

import io.confluent.kafka.serializers.KafkaAvroDeserializer

...

val valueDeserializer = new KafkaAvroDeserializer()
valueDeserializer.configure(
  Map(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG -> schemaRegistryUrl).asJava, 
  false)

...

override def deserialize(messageKey: Array[Byte], message: Array[Byte], 
                       topic: String, partition: Int, offset: Long): KafkaKV = {

    val key = keyDeserializer.deserialize(topic, messageKey).asInstanceOf[GenericRecord]
    val value = valueDeserializer.deserialize(topic, message).asInstanceOf[GenericRecord]

    KafkaKV(key, value)
    }

...

这里有详细解释:http://svend.kelesia.com/how-to-integrate-flink-with-confluents-schema-registry.html#how-to-integrate-flink-with-confluents-schema-registry

希望对你有帮助!

【讨论】:

  • 在上面的 Java 示例和 Scala sn-p 中,我仍然对如何使用我的 .avsc avro 模式文件感到困惑。我需要使用 avro 工具 jar 编译它吗?我做到了,它从单个模式文件中生成了一些不同的 java 代码,可能是因为 .avsc 模式中每个不同的嵌套结构都有一个单独的 java 代码文件。是否有关于如何使用“.avsc”文件进行反序列化的教程?我希望它像 - new Deserializer("msg.avsc") 一样简单。即使是 Confluent 模式注册表也没有单独向模式注册表提供模式的说明。
猜你喜欢
  • 2019-06-29
  • 2019-08-05
  • 1970-01-01
  • 2019-07-12
  • 1970-01-01
  • 2019-11-25
  • 1970-01-01
  • 1970-01-01
  • 2020-12-13
相关资源
最近更新 更多