【发布时间】:2019-09-18 13:47:00
【问题描述】:
我正在尝试创建一个通用实用程序,它将从 Kafka 主题读取 avro 文件并将 avro 文件写入 Java 中的主题。 我找不到太多相同的文档。 感谢任何有效的代码。
【问题讨论】:
-
看看您到目前为止尝试了什么以及遇到了什么问题会很有用。否则,问题有点模糊
标签: java serialization apache-kafka deserialization avro
我正在尝试创建一个通用实用程序,它将从 Kafka 主题读取 avro 文件并将 avro 文件写入 Java 中的主题。 我找不到太多相同的文档。 感谢任何有效的代码。
【问题讨论】:
标签: java serialization apache-kafka deserialization avro
也许你看到了这个问题? 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 接口也可以使用双射
【讨论】:
Unknown magic byte 表示 some 事件实际上不是 Avro。没有简单的方法来读取这些事件,因此您只需捕获异常并跳过它们或将它们发送到不同的“死信”主题
kafka-console-consumer --topic partiton_1 --partition 6 --offset 15047 --max-messages=1查看事件