【问题标题】:Avro Serialization/Deserialization to/from Kafka Topic [closed]Avro序列化/反序列化到/从Kafka主题[关闭]
【发布时间】:2019-09-18 13:47:00
【问题描述】:

我正在尝试创建一个通用实用程序,它将从 Kafka 主题读取 avro 文件并将 avro 文件写入 Java 中的主题。 我找不到太多相同的文档。 感谢任何有效的代码。

【问题讨论】:

  • 看看您到目前为止尝试了什么以及遇到了什么问题会很有用。否则,问题有点模糊

标签: java serialization apache-kafka deserialization avro


【解决方案1】:

也许你看到了这个问题? Read Existing Avro File and Send to Kafka


您通常在 Kafka 中没有“文件”...围绕 Avro 有很多关于如何读取/写入文件的文档,但 Kafka 仅将单个记录处理为 byte[] 对象。 Avro 提供了BinaryEncoder 类来获取字节数组的记录

如果您将 Kafka 与 Avro 结合使用,您通常会使用 Confluent Schema Registry。这使得每条 Kafka 消息都不需要完全编码的 Avro 模式,而只需一个带有二进制数据的数字参考 id

你可以在这里找到他们的快速入门

https://docs.confluent.io/current/quickstart/index.html

这里还有 Github 示例代码库

https://github.com/confluentinc/examples/blob/5.2.1-post/clients/avro/README.md


如果您不使用架构注册表,则必须编写自己的序列化程序。这是一个通过 Bijection 库为生产者使用普通 Kafka API 和为消费者使用 Spark 的示例

http://aseigneurin.github.io/2016/03/04/kafka-spark-avro-producing-and-consuming-avro-messages.html

请注意,Spark 已经有一个用于处理 Avro 的包。理论上,您可以直接使用它来读取 Avro 文件作为 Dataframe 并将它们写入 Kafka 主题。

Spark 完全没有必要。 Kafka Consumer 或 Deserializer 接口也可以使用双射

【讨论】:

  • 创建生产者和消费者类(使用 schemaregistry)。当我使用 io.confluent.kafka.serializers.KafkaAvroDeserializer 时,有时会抛出错误说 org.apache.kafka.common.errors.SerializationException:在偏移量 15047 处反序列化分区 partiton_1-6 的键/值时出错
  • Unknown magic byte 表示 some 事件实际上不是 Avro。没有简单的方法来读取这些事件,因此您只需捕获异常并跳过它们或将它们发送到不同的“死信”主题
  • 您可以在 CLI 上使用kafka-console-consumer --topic partiton_1 --partition 6 --offset 15047 --max-messages=1查看事件
猜你喜欢
  • 2018-08-11
  • 1970-01-01
  • 2019-07-30
  • 2015-08-01
  • 1970-01-01
  • 2017-11-20
  • 2019-11-18
  • 2017-01-22
  • 2019-02-12
相关资源
最近更新 更多