Flink 集成 Confluent Avro 格式Schema Registry 序列化/反序列化完整指南【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flinkavro-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 版本 idschema 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 iddataInputStream.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 byte0与 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 Registrybootstrap 地址为localhost:9092Schema 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_keyvalue 使用avro-confluent格式value.avro-confluent.url指向 Schema Registryvalue.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 valueCREATE 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-confluentavro-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 Tokenavro-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 格式时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 Registryvalue.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 Registryvalue.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.javaConfluent 协议编解码ConfluentSchemaRegistryCoder.java序列化/反序列化 SchemaConfluentRegistryAvroSerializationSchema.java、ConfluentRegistryAvroDeserializationSchema.javaDebezium CDC 变体debezium-avro-confluent支持完整 changelogDebeziumAvroFormatFactory.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),仅供参考
企业数字化 ERP 产品动态
相关推荐
3个坑让校内人人网变慢?实战项目性能优化全解 3个坑让校内人人网变慢?实战项目性能优化全解 面试被问“为什么列表加载慢”,你答不上来?别慌。 很多在校招或社招中,候选人死就死在 实战项目 的细节上。 特别是像【校内人人网】这种典型的B/S架构项目,性能瓶颈往往藏在不起眼的地方。… · 2026/9/23 4:56:03
STM32+CODESYS:低成本工业PLC开发实战全攻略 做工业控制这一行,选型永远是第一道坎。前年我接了一个小型包装线控制系统的改造项目,甲方预算压得很死,还要求支持以太网远程监控和在线修改参数。传统中型PLC加组态软件授权的方案,光软件和CPU模块就超了预算的一半,… · 2026/9/23 4:56:03
栈数据结构深度解析:从函数调用到堆栈溢出实战 1. 堆栈究竟是什么:从生活场景到核心抽象如果你接触过数据结构,哪怕只是刚开始准备考研、刷LeetCode,或者在学校里正在为《数据结构》实验报告发愁,那“堆栈”这个词你一定不陌生。它还有个别名叫“栈”,英文叫Stack。… · 2026/9/23 5:39:44
management缩写避坑指南:3个常见误区+完整示例 management缩写避坑指南:3个常见误区+完整示例 官方文档翻了三遍还是记不住 management 的缩写?别慌,这不是你笨,是文档写法反人类。我见过太多开发者在配置 API 或解析日志时,因为搞混 mgmt 、 mgt 、… · 2026/9/23 5:39:44
SpringBoot+Vue箱包仓储管理系统全栈开发实践 1. 项目概述:箱包存储系统信息管理解决方案箱包存储系统信息管理系统是一套针对仓储物流行业设计的全栈解决方案,它完美结合了SpringBoot后端的高效稳定、Vue前端的灵活交互以及MySQL的数据可靠性。这个开箱即用的系统特别适合中小型物流企业、电商仓库以… · 2026/9/23 5:39:44
搞定n代表什么数:附完整示例与性能优化实战 搞定n代表什么数:附完整示例与性能优化实战 你复制来的代码跑不通,是不是经常卡在这里?别急,今天我们不聊虚的,直接上 完整示例 ,带你彻底搞懂循环变量 n 在性能优化里的坑。很多老手都栽在这上面,以为 n 只是个数,其实它决定了你的算法是… · 2026/9/23 5:39:44
从拟声词到网络热词:cua如何成为全网通用梗 最近不管是刷短视频,还是混在各种聊天群里,你大概率都见过这个词:cua。别看它只有三个字母,现在它已经不是一个简单的拟声词了,而是一种状态、一种情绪、一种“事情就这么发生了”的叙述方式。今天想好好聊聊它&#x… · 2026/9/23 5:39:37
QEMU ACPI PCI 热插拔接口规范:从 IO 端口协议到 AML 生成的完整解析 QEMU ACPI PCI 热插拔接口规范:从 IO 端口协议到 AML 生成的完整解析 【免费下载链接】qemu Official QEMU mirror. Please see https://www.qemu.org/contribute/ for how to submit changes to QEMU. Pull Requests are disabled. Please only use release tarbal… · 2026/9/23 5:39:37
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29