【问题标题】:How to use Avro generated schemas with Kafka Connect custom source connector如何将 Avro 生成的模式与 Kafka Connect 自定义源连接器一起使用
【发布时间】:2020-06-27 15:59:46
【问题描述】:

我正在使用 Kafka Connect 开发自定义源连接器,并且我正在尝试合并 Avro 支持。为此,我创建了一些 .avsc 文件来描述我的键和值模式,并将 avro-maven-plugin 添加到我的项目中,以自动创建可以在我的代码中使用的相应 java 类。

从生成的类中,我可以将架构作为org.apache.avro.Schema 类型的对象获取。

但是源连接器的poll 方法的返回类型是org.apache.kafka.connect.source.SourceRecord 对象的列表,其构造函数将模式作为org.apache.kafka.connect.data.Schema 的实例,我根本看不到直接的方法将一个转换为另一个。

那么我如何获得合适的键/值模式实例,然后我可以将它们插入到 SourceRecords 中,以便从连接器中的 poll 方法返回?

我在使用 Avro Maven 插件方面是否走在正确的轨道上,还是应该使用其他东西?

【问题讨论】:

  • 欢迎来到 StackOverflow!我只熟悉 Kafka Connect 框架的整体,而不是编写连接器的细节——但是如果你正在编写一个连接器,那么你不需要对 Avro 做任何事情,因为这是由 处理的转换器 下游。这些资源中的任何一个有帮助吗? opencredo.com/blogs/…docs.confluent.io/current/connect/devguide.html
  • 我已经阅读了您链接的页面,它们对我没有帮助。问题是我希望架构定义(即 .avsc 文件)独立于我的项目而存在。这应该是微服务架构的一部分,并且模式是一种 API,服务用于相互通信。如果我让 Kafka Connect 从我的代码中生成和注册模式,那么它们将不容易被其他甚至可能不是用 Java 编写的服务所使用,所以我想以某种方式反转该过程并使代码遵循模式。
  • 然后你需要为你的模式创建一个单独的项目,并将它们上传到一个 Maven 服务以供其他项目使用。

标签: apache-kafka avro apache-kafka-connect


【解决方案1】:

我不确定是否建议这样做,但是,您可以利用kafka-connect-avro-converter 库中提供的AvroData 类来完成转换。

图书馆可以在这里找到: https://mvnrepository.com/artifact/io.confluent/kafka-connect-avro-converter/5.4.1

该课程的来源在这里: https://github.com/confluentinc/schema-registry/blob/5.4.1-post/avro-converter/src/main/java/io/confluent/connect/avro/AvroData.java

您必须实例化 AvroData,然后尝试 toConnectSchema 函数。

【讨论】:

  • 这确实有效,但不知何故感觉不对。我认为如果有一种方法可以避免使用 avro 插件并仅从 .avsc 文件生成 Connect 模式定义,那会更简单,但似乎没有用于此的工具。也许我必须自己按照这些思路构建一些东西才能完成这项工作。
  • 我接受了这个答案,因为从技术上讲它是有效的,所以它是 a 解决方案,即使不一定是 解决方案。我决定采用我之前提交中概述的方式:编写一个自定义 maven 插件,直接生成用于创建 Connect 模式的 java 代码。
【解决方案2】:

您不应该在 Kafka Connect 中需要 Avro Schema。

Kafka Connect 维护一个内部 StructSchema 类,您应该在 SourceRecord / SinkRecord 类之间传递。例如,HTTP 源可以在 Struct 类中定义 int:statusstring:body

基本上,让Converter 接口负责任何和所有序列化。

【讨论】:

  • 你能用另一种方式解释吗?我认为我们需要传递用于在属性“value.schema”中创建对象(以及用于序列化和反序列化)的 avro 模式:但似乎这会导致 Cannot deserialize value of type `org.apache.kafka.connect.data.Schema$Type` from String "record": not one of the values accepted for Enum class: [STRING, INT16, STRUCT, BOOLEAN, ARRAY, FLOAT64, BYTES, MAP, INT64, INT32, INT8, FLOAT32] 如果我不这样做'不提供它返回的架构```找不到任何输入文件来推断架构。``我可以重用我的 .avsc 架构吗?
  • @Georgi 我不知道你指的是什么。 AvroConverter 没有这种名为 value.schema 的属性,但我的回答是指自定义编写的连接器应该只使用 Schema 和 Struct 类来表示内部数据
  • @OneCrickleteer 我的意思是不要将 JSON 文件中的数据发送到几个 kafka 主题并使用模式注册表。我将值转换器设置为 AvroConverter。然后我将 auto.register.schemas 设置为 true 并期望当我将 json 提供给配置时,它应该单独生成模式并将其注册到 SR 中。不幸的是,这失败了,因为找不到任何输入文件来推断架构。然后我发现在 kafka-connect jsonspooldir 中有 value.schema 选项,但是如果我将我的 avro 模式传递给它(转义和缩小)它会失败并出现提到的字符串记录错误......我把事情搞砸了吗?
  • @Georgi 我建议为此创建自己的帖子。我从未使用过 spooldir 连接器,但我认为它使用的架构不是 Avro 语法
猜你喜欢
  • 2020-12-08
  • 2020-04-09
  • 2020-01-07
  • 1970-01-01
  • 2021-12-25
  • 2021-04-07
  • 2018-08-20
  • 2021-07-06
  • 2020-08-05
相关资源
最近更新 更多