【问题标题】:Can I use confluent Schema Registry to generate schema less avro msgs from flat file?我可以使用 confluent Schema Registry 从平面文件中生成 schemaless avro msgs 吗?
【发布时间】:2017-12-17 15:23:43
【问题描述】:

我想知道我可以使用 Confluent Schema 注册表来生成(然后将其发送到 kafka)架构少 avro 记录吗?如果是的话,有人可以分享一些资源吗? 我在 Confluent 网站和 Google 上找不到任何示例。

我有一个普通的分隔文件,我有一个单独的架构,目前我正在使用 Avro 通用记录架构来序列化 Avro 记录并通过 Kafka 发送它。这样,模式仍然与记录相关联,这使得它更加庞大。我的逻辑是,如果我在从 kafka 发送记录时删除模式,我将能够获得更高的吞吐量。

【问题讨论】:

  • 为什么要使用模式注册表来发送无模式记录?我很困惑。
  • 实际上我目前正在使用通用记录 Avro 模式从 csv 生成 Avro 记录,所以我的理解是它是在将模式发送到 kafka 时将模式附加到 Avro 二进制记录,这使得我的 Kafka 负载更大。
  • 我不知道您可以在本地将 Avro 与数据中包含的模式分离......但是,Kafka 似乎为 Avro 实现了特定的序列化程序并剥离了 Avro 模式以进行传输:@ 987654321@

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


【解决方案1】:

Confluent Schema Registry 将发送序列化的 Avro 消息,而消息中不包含整个 Avro Schema。我认为这就是您所说的“少模式”消息的意思。

Confluent Schema Registry 将存储 Avro 模式,并且在线消息中只包含一个简短的索引 ID。

这里有完整的文档,包括测试 Confluent Schema Registry 的快速入门指南

http://docs.confluent.io/current/schema-registry/docs/index.html

【讨论】:

  • 感谢您的回答,我有一个纯分隔文件,那么如何将注册表模式 ID 附加到它并通过 kafka 发送?你有什么例子吗?
  • 文档中有如何使用 Kafka Java Producer API 发布的示例。具体在这里docs.confluent.io/current/schema-registry/docs/…
  • 您还可以使用 REST API 将 Avro 消息发布到 Kafka。文档和示例在这里docs.confluent.io/current/kafka-rest/docs/intro.html
  • 谢谢汉斯,我正在检查序列化程序代码,我有一个问题,为什么即使架构已经在架构注册表中注册,还需要使用 userSchema?
  • 这是可选的,但发布者可能正在注册比注册表中的架构更新的版本。
【解决方案2】:

您可以在 cmd 的以下命令的帮助下首次注册您的 avro 架构

curl -X POST -i -H "Content-Type: application/vnd.schemaregistry.v1+json" \
        --data '{"schema": "{\"type\": \"string\"}"}' \
        http://localhost:8081/subjects/topic

您可以使用

查看主题的所有版本
curl -X GET -i http://localhost:8081/subjects/topic/versions

要从融合模式注册表中存在的所有版本中查看完整的 Acro 模式,请使用以下命令,将以 json 格式显示模式

  curl -X GET -i http://localhost:8081/subjects/topica/versions/1

Avro 模式注册是 Kafka 生产者的任务

在融合模式注册表中拥有模式后,您只需要将 avro 通用记录发布到特定的 kafka 主题,在我们的例子中是“主题”

Kafka Consumer:使用以下代码获取特定 Kafka 主题的最新架构

val schemaReg = new CachedSchemaRegistryClient(kafkaAvroSchemaRegistryUrl, 100)
val schemaMeta = schemaReg.getLatestSchemaMetadata(kafkaTopic + "-value")
val schema = schemaMeta.getSchema
val schema =new Schema.Parser().parse(schema)

上面将用于获取模式,然后我们可以使用 confluent 解码来自 kafka 主题的记录。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-10-06
    • 1970-01-01
    • 1970-01-01
    • 2018-05-17
    • 1970-01-01
    • 1970-01-01
    • 2022-12-27
    • 2021-10-04
    相关资源
    最近更新 更多