【问题标题】:Exception while Deserialize avro data using ConfluentSchemaRegistry?使用 ConfluentSchemaRegistry 反序列化 avro 数据时出现异常?
【发布时间】:2019-09-11 18:10:05
【问题描述】:

我是 flink 和 Kafka 的新手。我正在尝试使用 Confluent Schema 注册表反序列化 avro 数据。我已经在 ec2 机器上安装了 flink 和 Kafka。此外,在运行代码之前已经创建了“测试”主题。

代码路径:https://gist.github.com/mandar2174/5dc13350b296abf127b92d0697c320f2

代码执行以下操作作为实现的一部分:

1) Create a flink DataStream object using a list of user element. (User class is avro generated class)
2) Write the Datastream source to Kafka using AvroSerializationSchema.
3) Read the data from Kafka using ConfluentRegistryAvroDeserializationSchema by reading the schema from Confluent Schema registry.

运行 flink 可执行 jar 的命令:

./bin/flink run -c com.streaming.example.ConfluentSchemaRegistryExample /opt/flink-1.7.2/kafka-flink-stream-processing-assembly-0.1.jar

运行代码时出现异常:

java.io.IOException: Unknown data format. Magic number does not match
    at org.apache.flink.formats.avro.registry.confluent.ConfluentSchemaRegistryCoder.readSchema(ConfluentSchemaRegistryCoder.java:55)
    at org.apache.flink.formats.avro.RegistryAvroDeserializationSchema.deserialize(RegistryAvroDeserializationSchema.java:66)
    at org.apache.flink.streaming.util.serialization.KeyedDeserializationSchemaWrapper.deserialize(KeyedDeserializationSchemaWrapper.java:44)
    at org.apache.flink.streaming.connectors.kafka.internal.KafkaFetcher.runFetchLoop(KafkaFetcher.java:140)
    at org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase.run(FlinkKafkaConsumerBase.java:665)
    at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:94)
    at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:58)
    at org.apache.flink.streaming.runtime.tasks.SourceStreamTask.run(SourceStreamTask.java:99)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:300)
    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:704)
    at java.lang.Thread.run(Thread.java:748)

我用于 User 类的 Avro 架构如下:

{
  "type": "record",
  "name": "User",
  "namespace": "com.streaming.example",
  "fields": [
    {
      "name": "name",
      "type": "string"
    },
    {
      "name": "favorite_number",
      "type": [
        "int",
        "null"
      ]
    },
    {
      "name": "favorite_color",
      "type": [
        "string",
        "null"
      ]
    }
  ]
}

有人能指出我在使用融合 Kafka 模式注册表反序列化 avro 数据时缺少哪些步骤吗?

【问题讨论】:

    标签: apache-kafka apache-flink avro confluent-platform confluent-schema-registry


    【解决方案1】:

    您编写 Avro 数据的方式也需要使用注册表,以便依赖于它的反序列化器工作。

    But this is an open PR in Flink, still 用于添加 ConfluentRegistryAvroSerializationSchema

    我相信解决方法是使用AvroDeserializationSchema,它不依赖于注册表。

    如果您确实想在生产者代码中使用 Registry,那么您必须在 Flink 之外这样做,直到该 PR 被合并。

    【讨论】:

    • 我的理解是否正确“我不能直接将 AvroSerializationSchema 与 ConfluentRegistryAvroDeserializationSchema 一起使用,因为序列化和反序列化都应该参考融合模式注册表”?你有任何参考示例代码,我可以参考在生产者代码中使用模式注册表吗?
    • 对。据我所知,Flink Avro 序列化程序不使用注册表,因此 Avro 模式将作为每条消息的一部分嵌入。可以在此处查看纯 Java 生产者示例docs.confluent.io/current/schema-registry/…
    • 感谢您的确认。我将尝试参考示例代码并使用融合模式注册表测试序列化和反序列化。测试后我会更新我的结果。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-09
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多