【问题标题】:How to utilize existing avro schema for my kafka consumer?如何为我的 kafka 消费者使用现有的 avro 模式?
【发布时间】:2021-05-09 04:26:32
【问题描述】:

我正在使用 Debezium SQL Server 连接器进行更改数据捕获,连接器会自动生成架构并将架构注册到架构注册表中,这意味着我没有 avro 架构文件。在这种情况下,如何使用此模式编写消费者读取数据?我看过很多文章使用 avro schema 文件为消费者读取数据,schema 注册表中只会有这个payload的一个schema。

如果我在本地创建一个 avro 文件并让我的消费者使用它,那么我必须使用不同的名称注册一个重复的架构。

我的问题是如何使用 kafka 连接器注册的这个模式编写 Java 消费者 API。非常感谢。

这是我的价值模式:

{"subject":"new.dbo.locations-value","version":1,"id":102,"schema":"{\"type\":\"record\",\"name\":\"Envelope\",\"namespace\":\"new.dbo.locations\",\"fields\":[{\"name\":\"before\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"Value\",\"fields\":[{\"name\":\"id\",\"type\":\"long\"},{\"name\":\"display_id\",\"type\":\"string\"},{\"name\":\"first_name\",\"type\":\"string\"},{\"name\":\"last_name\",\"type\":\"string\"},{\"name\":\"location_id\",\"type\":\"string\"},{\"name\":\"location_name\",\"type\":\"string\"},{\"name\":\"location_sub_type_id\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"location_time_zone\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"parent_organization_id\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"type\",\"type\":[\"null\",\"string\"],\"default\":null}],\"connect.name\":\"new.dbo.locations.Value\"}],\"default\":null},{\"name\":\"after\",\"type\":[\"null\",\"Value\"],\"default\":null},{\"name\":\"source\",\"type\":{\"type\":\"record\",\"name\":\"Source\",\"namespace\":\"io.debezium.connector.sqlserver\",\"fields\":[{\"name\":\"version\",\"type\":\"string\"},{\"name\":\"connector\",\"type\":\"string\"},{\"name\":\"name\",\"type\":\"string\"},{\"name\":\"ts_ms\",\"type\":\"long\"},{\"name\":\"snapshot\",\"type\":[{\"type\":\"string\",\"connect.version\":1,\"connect.parameters\":{\"allowed\":\"true,last,false\"},\"connect.default\":\"false\",\"connect.name\":\"io.debezium.data.Enum\"},\"null\"],\"default\":\"false\"},{\"name\":\"db\",\"type\":\"string\"},{\"name\":\"schema\",\"type\":\"string\"},{\"name\":\"table\",\"type\":\"string\"},{\"name\":\"change_lsn\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"commit_lsn\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"event_serial_no\",\"type\":[\"null\",\"long\"],\"default\":null}],\"connect.name\":\"io.debezium.connector.sqlserver.Source\"}},{\"name\":\"op\",\"type\":\"string\"},{\"name\":\"ts_ms\",\"type\":[\"null\",\"long\"],\"default\":null},{\"name\":\"transaction\",\"type\":[\"null\",{\"type\":\"record\",\"name\":\"ConnectDefault\",\"namespace\":\"io.confluent.connect.avro\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"total_order\",\"type\":\"long\"},{\"name\":\"data_collection_order\",\"type\":\"long\"}]}],\"default\":null}],\"connect.name\":\"new.dbo.locations.Envelope\"}"}%

【问题讨论】:

    标签: java apache-kafka kafka-consumer-api avro confluent-schema-registry


    【解决方案1】:

    您不需要本地架构文件。您可以使用KafkaConsumer<?, GenericRecord> 进行消费,这将使反序列化程序下载并缓存每条消息的相应 ID+schema。

    这种方法的缺点是您需要小心解析数据(很像原始 JSON)

    如果您需要静态架构和允许严格类型检查的编译类,请从/subjects/:name/versions/latest 的注册表下载它

    【讨论】:

    • 非常感谢!这个对我有用!顺便说一句,我如何使用 KafkaConsumer, GenericRecord> 为我的消费者 API 选择特定版本的架构?
    • 如果您需要特定版本,即您确实需要使用/下载架构文件。使用这种方法,它总是反序列化为嵌入在消息中的模式 ID 版本(不幸的是,没有办法在运行时实际检测到该版本)
    • 感谢您的快速回复。很有帮助!
    猜你喜欢
    • 2018-07-04
    • 2021-03-06
    • 2016-12-09
    • 2014-07-30
    • 2023-01-10
    • 2020-08-20
    • 1970-01-01
    • 2021-05-14
    • 1970-01-01
    相关资源
    最近更新 更多