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.format或value.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 协议:
- 读 schema:
readSchema()首先读取一个 magic byte 并校验其值必须为0(CONFLUENT_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"。 - 写 schema:
writeSchema()先调用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",同时实现DeserializationFormatFactory与SerializationFormatFactory:
- 反序列化侧(
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.confluent、org.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.url与value.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) | String | Schema Registry 的 Basic Auth 凭据来源 |
avro-confluent.basic-auth.user-info | 可选 | (none) | String | Schema Registry 的 Basic Auth 用户信息 |
avro-confluent.bearer-auth.credentials-source | 可选 | (none) | String | Schema Registry 的 Bearer Auth 凭据来源 |
avro-confluent.bearer-auth.token | 可选 | (none) | String | Schema Registry 的 Bearer Auth Token |
avro-confluent.properties | 可选 | (none) | Map | 透传给底层 Schema Registry 客户端的属性映射,适用于 Flink 尚未官方暴露的选项;注意 Flink 显式选项优先级更高 |
avro-confluent.ssl.keystore.location | 可选 | (none) | String | SSL keystore 的位置/文件 |
avro-confluent.ssl.keystore.password | 可选 | (none) | String | SSL keystore 的密码 |
avro-confluent.ssl.truststore.location | 可选 | (none) | String | SSL truststore 的位置/文件 |
avro-confluent.ssl.truststore.password | 可选 | (none) | String | SSL 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 格式时,kafka与upsert-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必填——这正是工厂在createEncodingFormat中subject.isPresent()校验的来由。SSL 与鉴权属性的透传:
buildOptionalPropertiesMap()(RegistryAvroFormatFactory.java)会把 Flink 显式维护的选项映射为 Schema Registry 客户端原生属性键,例如:avro-confluent.ssl.keystore.location→schema.registry.ssl.keystore.locationavro-confluent.basic-auth.user-info→basic.auth.user.infoavro-confluent.bearer-auth.token→bearer.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_INFO、URL等),具体取值含义以 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),仅供参考