【发布时间】: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