Avro序列化
Avro 序列化
Debezium 连接器在 Kafka Connect 框架中运行,通过生成变更事件记录来捕获数据库中的每一行级变更。对于每条变更事件记录,Debezium 连接器会完成以下操作:
- 应用已配置的转换(transformation)。
- 使用已配置的 Kafka Connect 转换器将记录的键和值序列化为二进制形式。
- 将记录写入正确的 Kafka 主题。
你可以为每个 Debezium 连接器实例分别指定转换器。Kafka Connect 提供了一个 JSON 转换器,可将记录的键和值序列化为 JSON 文档。默认情况下,JSON 转换器会包含记录的消息 schema,这会使每条记录的内容非常冗长。Debezium 教程展示了同时包含负载(payload)和 schema 时记录的样子。如果你希望以 JSON 序列化记录,可以考虑将以下连接器配置属性设置为 false:
key.converter.schemas.enablevalue.converter.schemas.enable
将这些属性设置为 false 后,每条记录中就不会再包含冗长的 schema 信息。
或者,你也可以使用 Apache Avro 来序列化记录的键和值。Avro 的二进制格式紧凑且高效。Avro schema 可以确保每条记录都具有正确的结构。Avro 的 schema 演进机制支持 schema 的演进。这对于 Debezium 连接器至关重要,因为连接器会动态生成每条记录的 schema,使其与被更改的数据库表结构相匹配。随着时间的推移,写入同一 Kafka 主题的变更事件记录可能会包含同一 schema 的不同版本。Avro 序列化使变更事件记录的消费者更容易适应不断变化的记录 schema。
要使用 Apache Avro 序列化,你必须部署一个用于管理 Avro 消息 schema 及其版本的 schema 注册表(schema registry)。可用的方案包括 Apicurio API and Schema Registry 以及 Confluent Schema Registry。下面对两者都做了介绍。
Apicurio API and Schema Registry
关于 Apicurio API and Schema Registry
Apicurio Registry 开源项目提供了多个可与 Avro 配合使用的组件:
一个可在 Debezium 连接器配置中指定的 Avro 转换器。该转换器将 Kafka Connect schema 映射为 Avro schema,然后使用这些 Avro schema 将记录的键和值序列化为 Avro 的紧凑二进制形式。
一个 API 和 schema 注册表,用于跟踪:
- Kafka 主题中使用的 Avro schema。
- Avro 转换器将生成的 Avro schema 发送到何处。
由于 Avro 模式存储在该注册表中,每条记录只需包含一个很小的 schema 标识符。这使每条记录变得更小。对于 Kafka 这样的 I/O 密集型系统,这意味着生产者和消费者可以获得更高的总吞吐量。
- 用于 Kafka 生产者和消费者的 Avro Serdes(序列化器和反序列化器)。你编写的用于消费变更事件记录的 Kafka 消费者应用程序,可以使用 Avro Serdes 对变更事件记录进行反序列化。
- 通过设置
as-confluent: true启用 Apicurio 的 Confluent 兼容模式,使 Apicurio Registry 能够使用与 Confluent Schema Registry 相同的线上格式(wire format)序列化 Kafka 消息(包括魔数字节和 4 字节的 schema ID)。这使得 Apicurio 能够与 Confluent 客户端和工具完全互操作,成为无需更改生产者或消费者的即插即用替代方案。
要在 Debezium 中使用 Apicurio Registry,请将 Apicurio Registry 转换器及其依赖项添加到你用于运行 Debezium 连接器的 Kafka Connect 容器镜像中。
| Apicurio Registry 项目还提供了一个 JSON 转换器。该转换器结合了消息更简洁和 JSON 易于阅读的优点。消息本身不包含模式信息,只包含一个 schema ID。 |
|---|
要使用 Apicurio Registry 提供的转换器,你需要提供 apicurio.registry.url。 |
|---|
Apicurio Registry 部署概述
要部署使用 Avro 序列化的 Debezium 连接器,你必须完成三项主要任务:
部署一个 Apicurio API 和 Schema Registry 实例。
从安装包中将 Avro 转换器安装到插件目录中。如果你使用的是 Debezium Connect 容器镜像,则无需安装该软件包。有关更多信息,请参阅使用 Debezium 容器部署 Apicurio Registry。
通过如下设置配置属性,配置 Debezium 连接器实例以使用 Avro 序列化:
key.converter=io.apicurio.registry.utils.converter.AvroConverter key.converter.apicurio.registry.url=http://apicurio:8080/apis/registry/v2 key.converter.apicurio.registry.auto-register=true key.converter.apicurio.registry.find-latest=true value.converter=io.apicurio.registry.utils.converter.AvroConverter value.converter.apicurio.registry.url=http://apicurio:8080/apis/registry/v2 value.converter.apicurio.registry.auto-register=true value.converter.apicurio.registry.find-latest=true schema.name.adjustment.mode=avro
或者,你也可以按照如下方式使用 Confluent 兼容性模式:
key.converter=io.apicurio.registry.utils.converter.AvroConverter
key.converter.apicurio.registry.url=http://apicurio:8080/apis/registry/v2
key.converter.apicurio.registry.auto-register=true
key.converter.apicurio.registry.find-latest=true
key.converter.schemas.enable": "false"
key.converter.apicurio.registry.headers.enabled": "false"
key.converter.apicurio.registry.as-confluent": "true"
key.converter.apicurio.use-id: "contentId"
value.converter=io.apicurio.registry.utils.converter.AvroConverter
value.converter.apicurio.registry.url=http://apicurio:8080/apis/registry/v2
value.converter.apicurio.registry.auto-register=true
value.converter.apicurio.registry.find-latest=true
value.converter.schemas.enable": "false"
value.converter.apicurio.registry.headers.enabled": "false"
value.converter.apicurio.registry.as-confluent": "true"
value.converter.apicurio.use-id: "contentId"
schema.name.adjustment.mode=avro在 Apicurio Registry 的 Confluent 兼容模式下,参数的配置方式是专门设计的,无论在线格式上还是客户端预期上,都力求与 Confluent Schema Registry 的工作方式完全一致。
converter.apicurio.registry.as-confluent: true:强制 Apicurio 以与 Confluent 相同的结构(魔术字节 + 4 字节的 schema ID)序列化 Kafka 消息。converter.apicurio.use-id: true(或globalId):确保消息中只写入数字形式的 schema ID,与 Confluent 消费者的预期相匹配。converter.schemas.enable: false:防止将完整的 Avro schema 嵌入 Kafka 消息中,因为 Confluent 反序列化器并不期望这样(它们会根据 ID 从注册中心获取 schema)。converter.apicurio.registry.headers.enabled: false:禁用 Confluent 客户端无法识别的 Apicurio 专属元数据头,使消息格式保持简洁并兼容。
在内部,Kafka Connect 始终使用 JSON 键/值转换器来存储配置和偏移量。
使用 Debezium 容器部署 Apicurio Registry
在你的环境中,你可能希望使用 Debezium 提供的容器镜像来部署使用 Avro 序列化的 Debezium 连接器。请按照以下步骤操作。在此过程中,你将在 Debezium Kafka Connect 容器镜像上启用 Apicurio 转换器,并将 Debezium 连接器配置为使用 Avro 转换器。
前置条件
- 你已安装 Docker,并拥有创建和管理容器的足够权限。
- 你已下载要以 Avro 序列化方式部署的 Debezium 连接器插件。
操作步骤
部署一个 Apicurio Registry 实例。
以下示例使用一个非生产环境的内存模式 Apicurio Registry 实例:
docker run -it --rm --name apicurio \ -p 8080:8080 apicurio/apicurio-registry:3.2.5运行 Debezium 的 Kafka Connect 容器镜像,通过设置环境变量
ENABLE_APICURIO_CONVERTERS=true来启用 Apicurio,从而为其配置 Avro 转换器:docker run -it --rm --name connect \ --link kafka:kafka \ --link mysql:mysql \ --link apicurio:apicurio \ -e ENABLE_APICURIO_CONVERTERS=true \ -e GROUP_ID=1 \ -e CONFIG_STORAGE_TOPIC=my_connect_configs \ -e OFFSET_STORAGE_TOPIC=my_connect_offsets \ -e KEY_CONVERTER=io.apicurio.registry.utils.converter.AvroConverter \ -e VALUE_CONVERTER=io.apicurio.registry.utils.converter.AvroConverter \ -e CONNECT_KEY_CONVERTER=io.apicurio.registry.utils.converter.AvroConverter \ -e CONNECT_KEY_CONVERTER_APICURIO.REGISTRY_URL=http://apicurio:8080/apis/registry/v2 \ -e CONNECT_KEY_CONVERTER_APICURIO_REGISTRY_AUTO-REGISTER=true \ -e CONNECT_KEY_CONVERTER_APICURIO_REGISTRY_FIND-LATEST=true \ -e CONNECT_VALUE_CONVERTER=io.apicurio.registry.utils.converter.AvroConverter \ -e CONNECT_VALUE_CONVERTER_APICURIO_REGISTRY_URL=http://apicurio:8080/apis/registry/v2 \ -e CONNECT_VALUE_CONVERTER_APICURIO_REGISTRY_AUTO-REGISTER=true \ -e CONNECT_VALUE_CONVERTER_APICURIO_REGISTRY_FIND-LATEST=true \ -e CONNECT_SCHEMA_NAME_ADJUSTMENT_MODE=avro \ -p 8083:8083 quay.io/debezium/connect:3.6
Confluent Schema Registry
Confluent 提供了另一种 模式注册表 实现。
Confluent Schema Registry 部署概览
有关安装独立版 Confluent Schema Registry 的信息,请参阅 Confluent Platform 部署文档。
作为替代方案,你也可以将独立版 Confluent Schema Registry 作为容器进行 安装。
使用 Debezium 容器部署 Confluent Schema Registry
从 Debezium 2.0.0 开始,Debezium 容器不再包含对 Confluent Schema Registry 的支持。要为 Debezium 容器启用 Confluent Schema Registry,请将以下 Confluent Avro 转换器 JAR 文件安装到 Connect 插件目录中:
kafka-connect-avro-converterkafka-connect-avro-datakafka-avro-serializerkafka-schema-serializerkafka-schema-converterkafka-schema-registry-clientcommon-configcommon-utils
你可以从 Confluent Maven 仓库 下载上述文件。
还需要将其他一些 JAR 文件放到 Connect 插件目录中:
avrocommons-compressfailureaccessguavaminimal-jsonre2jslf4j-apisnakeyamlswagger-annotationsjackson-databindjackson-corejackson-annotationsjackson-dataformat-csvlogredactorlogredactor-metrics
你可以从 Maven 仓库 下载上述文件。
配置方式略有不同。
在 Debezium 连接器配置中指定以下属性:
key.converter=io.confluent.connect.avro.AvroConverter key.converter.schema.registry.url=http://localhost:8081 value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://localhost:8081部署一个 Confluent Schema Registry 实例:
docker run -it --rm --name schema-registry \ --link kafka \ -e SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS=kafka:9092 -e SCHEMA_REGISTRY_HOST_NAME=schema-registry \ -e SCHEMA_REGISTRY_LISTENERS=http://schema-registry:8081 \ -p 8181:8181 confluentinc/cp-schema-registry运行配置为使用 Avro 的 Kafka Connect 镜像:
运行一个控制台消费者,从
db.myschema.mytable主题读取新的 Avro 消息并解码为 JSON:docker run -it --rm --name avro-consumer \ --link kafka:kafka \ --link mysql:mysql \ --link schema-registry:schema-registry \ quay.io/debezium/connect:3.6 \ /kafka/bin/kafka-console-consumer.sh \ --bootstrap-server kafka:9092 \ --property print.key=true \ --formatter io.confluent.kafka.formatter.AvroMessageFormatter \ --property schema.registry.url=http://schema-registry:8081 \ --topic db.myschema.mytable
命名
正如 Avro 文档中所述,名称必须遵循以下规则:
- 以
[A-Za-z_]开头 - 其后仅包含
[A-Za-z0-9_]字符
Debezium 以列名作为对应 Avro 字段名称的基础。如果列名不符合 Avro 的命名规则,这可能会在序列化过程中导致问题。每个 Debezium 连接器都提供了一个配置属性 field.name.adjustment.mode,如果你的列名不符合 Avro 的命名规则,可以将该属性设置为 avro。将 field.name.adjustment.mode 设置为 avro 后,无需实际修改模式即可对不符合规范的字段进行序列化。
获取更多信息
Debezium 博客上的这篇文章介绍了序列化器、转换器及其他组件的概念,并讨论了使用 Avro 的优势。自该文章发布以来,部分 Kafka Connect 转换器的细节已略有变化。
有关将 Avro 用作 Debezium 变更数据事件消息格式的完整示例,请参阅 MySQL 与 Avro 消息格式。
评论
登录后参与评论
KnowForge