Flink 集成 Confluent Avro 格式:Schema Registry 序列化/反序列化完整指南
2026/9/23 4:56:21 网站建设 项目流程

Flink 集成 Confluent Avro 格式:Schema Registry 序列化/反序列化完整指南

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

avro-confluent是 Apache Flink 官方提供的一种序列化格式(Serialization Schema / Deserialization Schema),用于与 Confluent Schema Registry 协同工作:它可以读取由io.confluent.kafka.serializers.KafkaAvroSerializer序列化的记录,也可以写出能被io.confluent.kafka.serializers.KafkaAvroDeserializer反序列化的记录。本文以 avro-confluent.md 为骨架,结合本仓库 flink-avro-confluent-registry 模块源码,系统讲解格式工作原理、依赖引入方式、三类建表示例、全部可配置参数、鉴权/SSL 安全选项以及数据类型映射规则,读完即可在 Kafka / Upsert Kafka 表上落地 Avro + Schema Registry 的数据读写方案。

格式定位与工作原理

在 Flink Table / SQL 生态中,avro-confluent是一种序列化格式(format),它不独立存在,必须挂载在连接器之上。官方文档明确说明:该格式只能与 Apache Kafka SQL 连接器 或 Upsert Kafka SQL 连接器 一起使用,分别充当key.formatvalue.format

其核心工作模式分为读写两个方向:

  • 读取(反序列化):根据记录中编码的 schema 版本 id(schema id),从配置的 Confluent Schema Registry 中拉取 Avrowriter schema;而reader schema则由 Flink 的 table schema 推断而来。这保证了消费端可以兼容上游 Producer 写入时使用的历史 schema 版本。
  • 写入(序列化):Flink 从 table schema 推断出 Avro schema,将其注册到 Schema Registry 获取 schema id,并把 schema id 与数据一起编码进 Kafka 消息。schema 注册在哪个 subject 下,由avro-confluent.subject参数(或其 key/value 前缀变体)控制。

源码层的协议实现

底层协议的编解码逻辑位于 ConfluentSchemaRegistryCoder.java,它实现了 Flink 的SchemaCoder接口,完整复刻了 Confluent 的 wire format 协议:

  • 读 schemareadSchema()首先读取一个 magic byte 并校验其值必须为0CONFLUENT_MAGIC_BYTE),随后读取 4 字节的 schema id(dataInputStream.readInt()),最后通过schemaRegistryClient.getById(schemaId)从 Schema Registry 获取对应的 Avro schema。若 magic number 不匹配会抛出IOException("Unknown data format. Magic number does not match"),若查不到 schema 会提示 "Could not find schema with id ... in registry"。
  • 写 schemawriteSchema()先调用schemaRegistryClient.register(subject, schema)将 schema 注册到指定 subject 并获得 schema id,然后依次写出 magic byte(0)与 4 字节大端 schema id,与 Confluent 官方KafkaAvroSerializer的输出完全兼容。

该 Coder 通过 CachedSchemaCoderProvider.java 创建,内部使用CachedSchemaRegistryClient缓存 schema(默认 identity map 容量为 1000),避免每条记录都向 Schema Registry 发起 HTTP 请求;Schema Registry 客户端配置来自 URL 与额外属性映射。

Format Factory 的装配逻辑

RegistryAvroFormatFactory.java 定义了工厂标识符IDENTIFIER = "avro-confluent",同时实现DeserializationFormatFactorySerializationFormatFactory

  • 反序列化侧(createDecodingFormat)要求avro-confluent.url必填;若配置了avro-confluent.schema则校验其与 table schema 一致,否则通过AvroSchemaConverter.convertToSchema(rowType)由 table schema 直接转换出 reader schema。
  • 序列化侧(createEncodingFormat)同样要求 URL 必填,并且必须提供 subject,否则抛出ValidationException(测试 RegistryAvroFormatFactoryTest.java 中专门验证了缺失 subject 时的报错信息 "Option avro-confluent.subject is required for serialization")。
  • 注意,avro-confluent本身只支持insertOnly的 ChangelogMode;如需处理 Debezium 的 CDC 变更流(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE),可使用同模块下的debezium-avro-confluent格式(工厂见 DebeziumAvroFormatFactory.java)。

依赖引入

使用avro-confluent需要引入以下 SQL Jar(以本仓库2.0-SNAPSHOT版本为例):

  • flink-sql-avro-confluent-registry-2.0-SNAPSHOT.jar:面向 SQL 用户的打包产物。其 pom.xml 使用 maven-shade-plugin 将io.confluentorg.apache.kafka重定位为org.apache.flink.avro.registry.confluent.shaded.*,避免与用户 classpath 上的 Confluent/Kafka 依赖冲突;同时将org.apache.avro重定位到与flink-sql-avro相同的 shade 命名空间,从而允许两个 SQL Jar 同时存在。
  • 底层依赖模块:flink-avro-confluent-registry(核心实现,见 flink-avro-confluent-registry/pom.xml),其依赖io.confluent:kafka-schema-registry-client(本仓库锁定版本7.5.3)。

如果是 Maven、SBT、Gradle 等构建工具方式引入,还需要在构建文件中配置 Confluent 的 Maven 仓库https://packages.confluent.io/maven/,因为kafka-schema-registry-client发布在 Confluent 自己的仓库中(本仓库的 pom.xml 中同样声明了该<repository>)。

在 SQL Client 中引入 Jar 的典型做法是放在lib/目录下,或使用ADD JAR语句动态加载。

如何创建使用 avro-confluent 格式的表

以下三类示例分别覆盖 Kafka 连接器的 key/value 不同组合方式,以及 Upsert Kafka 连接器的用法,均取自官方文档并可直接复制运行(假设本机已有 Kafka 与 Confluent Schema Registry,bootstrap 地址为localhost:9092,Schema Registry 地址为localhost:8082)。

示例一:Kafka key 为 UTF-8 字符串,value 为 Schema Registry 中的 Avro 记录

CREATE TABLE user_created ( -- 该列映射到 Kafka 原始的 UTF-8 key the_kafka_key STRING, -- 映射到 Kafka value 中的 Avro 字段的一些列 id STRING, name STRING, email STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user_events_example1', 'properties.bootstrap.servers' = 'localhost:9092', -- UTF-8 字符串作为 Kafka 的 keys,使用表中的 'the_kafka_key' 列 'key.format' = 'raw', 'key.fields' = 'the_kafka_key', 'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'http://localhost:8082', 'value.fields-include' = 'EXCEPT_KEY' )

写入数据:

INSERT INTO user_created SELECT -- 将 user id 复制至映射到 kafka key 的列中 id as the_kafka_key, -- 所有的 values id, name, email FROM some_table

要点说明:

  • key 使用raw格式,key.fields指定哪一列承载 Kafka key(本例为the_kafka_key);
  • value 使用avro-confluent格式,value.avro-confluent.url指向 Schema Registry;
  • value.fields-include = 'EXCEPT_KEY'表示 value 的 Avro schema 只包含除 key 列之外的其余字段,避免 key 字段重复出现在 value 中。

示例二:Kafka 的 key 与 value 都注册为 Avro 记录

CREATE TABLE user_created ( -- 该列映射到 Kafka key 中的 Avro 字段 'id' kafka_key_id STRING, -- 映射到 Kafka value 中的 Avro 字段的一些列 id STRING, name STRING, email STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'user_events_example2', 'properties.bootstrap.servers' = 'localhost:9092', -- 注意:由于哈希分区,在 Kafka key 的上下文中,schema 升级几乎从不向后也不向前兼容。 'key.format' = 'avro-confluent', 'key.avro-confluent.url' = 'http://localhost:8082', 'key.fields' = 'kafka_key_id', -- 在本例中,我们希望 Kafka 的 key 和 value 的 Avro 类型都包含 'id' 字段 -- => 给表中与 Kafka key 字段关联的列添加一个前缀来避免冲突 'key.fields-prefix' = 'kafka_key_', 'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'http://localhost:8082', 'value.fields-include' = 'EXCEPT_KEY', -- 自 Flink 1.13 起,subjects 具有一个默认值, 但是可以被覆盖: 'key.avro-confluent.subject' = 'user_events_example2-key2', 'value.avro-confluent.subject' = 'user_events_example2-value2' )

要点说明:

  • key 与 value 均使用avro-confluent,分别配置key.avro-confluent.urlvalue.avro-confluent.url
  • 由于 key 的 Avro 类型(含id字段)与 value 的 Avro 类型(含id字段)字段名重叠,通过key.fields-prefix = 'kafka_key_'给与 key 关联的列加前缀,映射时 Flink 会剥离该前缀后与 Avro 字段匹配;
  • 文档特别提醒:因为 Kafka 按 key 哈希分区,key 上的 schema 升级几乎既不向后兼容也不向前兼容,变更 key schema 需要特别谨慎(例如避免直接修改 key 字段类型);
  • subject自 Flink 1.13 起有默认值(见下节),但可以被显式覆盖,本例即为自定义 subject 的示范。

示例三:Upsert Kafka 连接器 + Avro value

CREATE TABLE user_created ( -- 该列映射到 Kafka 原始的 UTF-8 key kafka_key_id STRING, -- 映射到 Kafka value 中的 Avro 字段的一些列 id STRING, name STRING, email STRING, -- upsert-kafka 连接器需要一个主键来定义 upsert 行为 PRIMARY KEY (kafka_key_id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'user_events_example3', 'properties.bootstrap.servers' = 'localhost:9092', -- UTF-8 字符串作为 Kafka 的 keys -- 在本例中我们不指定 'key.fields',因为它由表的主键决定 'key.format' = 'raw', -- 在本例中,我们希望 Kafka 的 key 和 value 的 Avro 类型都包含 'id' 字段 -- => 给表中与 Kafka key 字段关联的列添加一个前缀来避免冲突 'key.fields-prefix' = 'kafka_key_', 'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'http://localhost:8082', 'value.fields-include' = 'EXCEPT_KEY' )

要点说明:

  • upsert-kafka连接器要求声明PRIMARY KEY,且 key 字段由主键自动决定,无需再写key.fields
  • 由于主键列kafka_key_id在 Kafka key 与 value 中都会出现,同样使用key.fields-prefix规避字段冲突;
  • value 部分与 Kafka 连接器用法一致,仍由avro-confluent+ Schema Registry 管理 Avro schema。

Format 参数详解

下表完整列出avro-confluent格式的全部参数(在 SQL 中使用时需按所在位置加前缀:作为 value 格式时前缀为value.avro-confluent.,作为 key 格式时前缀为key.avro-confluent.;不带前缀的format参数本身固定为'avro-confluent')。这些参数的键名与 AvroConfluentFormatOptions.java 中定义的ConfigOption一一对应。

参数是否必选默认值类型描述
format必选(none)String指定使用的格式,此处必须为'avro-confluent'
avro-confluent.basic-auth.credentials-source可选(none)StringSchema Registry 的 Basic Auth 凭据来源
avro-confluent.basic-auth.user-info可选(none)StringSchema Registry 的 Basic Auth 用户信息
avro-confluent.bearer-auth.credentials-source可选(none)StringSchema Registry 的 Bearer Auth 凭据来源
avro-confluent.bearer-auth.token可选(none)StringSchema Registry 的 Bearer Auth Token
avro-confluent.properties可选(none)Map透传给底层 Schema Registry 客户端的属性映射,适用于 Flink 尚未官方暴露的选项;注意 Flink 显式选项优先级更高
avro-confluent.ssl.keystore.location可选(none)StringSSL keystore 的位置/文件
avro-confluent.ssl.keystore.password可选(none)StringSSL keystore 的密码
avro-confluent.ssl.truststore.location可选(none)StringSSL truststore 的位置/文件
avro-confluent.ssl.truststore.password可选(none)StringSSL truststore 的密码
avro-confluent.schema可选(none)String在 Confluent Schema Registry 中已注册或待注册的 Avro schema;若不提供,Flink 将 table schema 转换为 Avro schema;提供的 schema 必须与 table schema 匹配
avro-confluent.subject可选(none)String序列化期间注册该格式所用 schema 的 Confluent Schema Registry subject;默认情况下,作为 value 或 key 格式时,kafkaupsert-kafka连接器使用<topic_name>-value<topic_name>-key作为默认 subject 名;但对于其他连接器(如filesystem),作为 sink 时该选项必填
avro-confluent.url必选(none)String用于获取/注册 schema 的 Confluent Schema Registry 地址

参数实现的源码佐证

  • 必选校验RegistryAvroFormatFactory.requiredOptions()仅返回URL一个必选项,因此无论读还是写,avro-confluent.url都是硬性要求。

  • subject 的默认值与必填校验:subject 的默认值逻辑(<topic_name>-value/<topic_name>-key)由 Kafka / Upsert Kafka 连接器注入(连接器已知 topic 名);当不经过 Kafka 连接器(如 filesystem sink)时,连接器无法推导 subject,因此avro-confluent.subject必填——这正是工厂在createEncodingFormatsubject.isPresent()校验的来由。

  • SSL 与鉴权属性的透传buildOptionalPropertiesMap()(RegistryAvroFormatFactory.java)会把 Flink 显式维护的选项映射为 Schema Registry 客户端原生属性键,例如:

    • avro-confluent.ssl.keystore.locationschema.registry.ssl.keystore.location
    • avro-confluent.basic-auth.user-infobasic.auth.user.info
    • avro-confluent.bearer-auth.tokenbearer.auth.token

    再与avro-confluent.properties中的自定义属性合并,最终一起交给CachedSchemaRegistryClient。测试 RegistryAvroFormatFactoryTest.java 同时覆盖了 Flink 显式选项与properties映射两种方式,并验证了schema与 table schema 不一致时的报错路径。

常见配置组合示例

带 Basic Auth 的 Schema Registry:

'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'https://schema-registry.example.com', 'value.avro-confluent.basic-auth.credentials-source' = 'USER_INFO', 'value.avro-confluent.basic-auth.user-info' = 'username:password'

带 SSL 与自定义属性的 Schema Registry:

'value.format' = 'avro-confluent', 'value.avro-confluent.url' = 'https://schema-registry.example.com', 'value.avro-confluent.ssl.keystore.location' = '/path/to/keystore.jks', 'value.avro-confluent.ssl.keystore.password' = 'keystore-password', 'value.avro-confluent.ssl.truststore.location' = '/path/to/truststore.jks', 'value.avro-confluent.ssl.truststore.password' = 'truststore-password', -- 透传 Flink 未官方暴露的 Schema Registry 客户端选项 'value.avro-confluent.properties.max.schemas.per.subject' = '1000'

说明:basic-auth.credentials-source的取值遵循 Confluent Schema Registry 客户端的约定(例如USER_INFOURL等),具体取值含义以 Confluent 客户端文档为准;上述示例中的 SSL 密码参数在仓库测试中以123456等示例值出现,生产环境请通过配置管理妥善保管。

数据类型映射

avro-confluent格式本身不提供独立的数据类型映射表,而是完全复用 Apache Avro Format 中定义的 Flink 数据类型与 Avro 类型对应关系(见其中#data-type-mapping一节)。

当前实现的行为要点:

  • 反序列化期间的 Avro reader schema 与序列化期间的 Avro writer schema,均由 Flink从 table schema 推断得到;显式定义 Avro schema 目前暂不支持(即不提供类似独立 Avro 格式那样"以 Avro schema 为准"的模式)。若通过avro-confluent.schema提供了 Avro schema 字符串,RegistryAvroFormatFactory.java 会先用AvroSchemaConverter.convertToDataType反解成 Flink 逻辑类型并与 table schema 比对,不一致则抛错,从而保证显式 schema 与 table schema 严格一致。
  • 除了映射表中列出的类型,Flink 还支持读写**可为空(nullable)**的类型:nullable 的 Flink 类型会被映射为 Avro 的union(something, null),其中something是由该 Flink 类型转换出的 Avro 类型。这一约定与 Confluent/标准 Avro 生态对 nullable 字段的表示一致。

关于 Avro 各类型的详细语义,可参考 Avro Specification(官方规范文档)。

关联实现与扩展阅读

  • 核心格式工厂:RegistryAvroFormatFactory.java
  • 配置项定义:AvroConfluentFormatOptions.java
  • Confluent 协议编解码:ConfluentSchemaRegistryCoder.java
  • 序列化/反序列化 Schema:ConfluentRegistryAvroSerializationSchema.java、ConfluentRegistryAvroDeserializationSchema.java
  • Debezium CDC 变体(debezium-avro-confluent,支持完整 changelog):DebeziumAvroFormatFactory.java
  • 单元测试:ConfluentSchemaRegistryCoderTest.java、RegistryAvroFormatFactoryTest.java
  • 普通 Avro 格式与数据类型映射:avro.md

需要处理 Debezium 产生的 CDC 数据时,可进一步阅读 debezium.md(其中debezium-avro-confluent与本文格式共享同一套 Schema Registry 编解码与安全配置体系)。

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询