
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Apache Flink 的 Avro format 允许用户基于 Avro schema 读取和写入 Avro 数据是 Kafka、Filesystem 等连接器与 Avro 序列化体系之间的桥梁。本文以 Flink 官方文档《Avro Format》为骨架结合仓库中flink-formats/flink-avro模块的源码实现完整讲解如何通过CREATE TABLE声明使用 Avro format、全部 format 参数的语义与默认值、Flink SQL 类型到 Avro 类型的映射规则以及 schema 推导与序列化/反序列化的底层工作原理。读完本文你将能够独立在 SQL 作业和 DataStream 作业中正确配置 Avro format并理解时间戳映射等易踩坑的兼容性细节。本文面向 Table SQL 生态对应中文文档为 docs/content.zh/docs/connectors/table/formats/avro.md其核心实现位于 flink-formats/flink-avro 模块。Avro Format 是什么Avro format 在 Flink 的 format 体系中同时扮演两种角色对应文档中的两个标签序列化 SchemaSerialization Schema与反序列化 SchemaDeserialization Schema。它允许基于 Avro schema 读取和写入 Avro 数据。目前Avro schema 不是由用户显式声明的而是从 table schema 推导而来。从源码看这一能力由 AvroFormatFactory.java 提供。它同时实现了DeserializationFormatFactory和SerializationFormatFactory工厂标识符为avro并在运行时分别构造AvroRowDataDeserializationSchemaAvro 字节 →RowData与AvroRowDataSerializationSchemaRowData→ Avro 字节两种 format 的ChangelogMode均为insertOnly()解码路径AvroRowDataDeserializationSchema内部先借助AvroDeserializationSchema.forGeneric(...)将消息反序列化为 AvroGenericRecord再通过AvroToRowDataConverters转换为 Flink 的RowData编码路径AvroRowDataSerializationSchema先通过RowDataToAvroConverters将RowData转换为GenericRecord再交给AvroSerializationSchema序列化为字节。引入依赖在 SQL Client / SQL Gateway 场景下Avro format 通过flink-sql-avro连接器 jar 提供。仓库中的 docs/data/sql_connectors.yml 定义了其坐标maven 制品名为flink-avroSQL jar 为flink-sql-avro-version.jar需要将其下载并放入 Flink 发行包的lib/目录后重启集群。在 Maven 项目中则直接依赖flink-avro模块即可模块定义见 flink-formats/flink-avro/pom.xmldependency groupIdorg.apache.flink/groupId artifactIdflink-avro/artifactId version${flink.version}/version /dependency该 pom 中org.apache.avro:avro版本由 Flink 父 pom 统一管理。另需注意pom 中joda-time被声明为providedoptional当 Avro 记录中含有 JodaTime 字段逻辑类型场景时使用者需要自行补充该依赖。使用 Avro format 创建表这是官方文档给出的使用 Kafka 连接器与 Avro format 创建表的完整示例CREATE TABLE user_behavior ( user_id BIGINT, item_id BIGINT, category_id BIGINT, behavior STRING, ts TIMESTAMP(3) ) WITH ( connector kafka, topic user_behavior, properties.bootstrap.servers localhost:9092, properties.group.id testGroup, format avro )要点解读format avro是 format 类型的声明入口Flink 的工厂查找机制会根据该值找到 AvroFormatFactory.java其IDENTIFIER常量正是avro该表声明了 5 个字段及其类型Avro schema 会据此自动推导Kafka 消息的 value 部分按 Avro 二进制编码解析表字段顺序与 Avro 记录字段一一对应ts TIMESTAMP(3)会按映射规则落到 Avro 的longtimestamp-millis详见下文类型映射与时间戳兼容性小节。Format 参数详解下表完整列出 Avro format 支持的全部参数继承自官方文档并结合 AvroFormatOptions.java 与 AvroFileFormatFactory.java 的源码补充了默认值与底层行为参数是否必选默认值类型描述format必选(none)String指定使用的 format这里应为avro。avro.codec可选snappyString仅用于 filesystem 连接器指定 Avro 的压缩编解码器。当前支持null、deflate、snappy、bzip2、xz。avro.encoding可选binaryString序列化采用的编码方式合法取值为binary与json。binary编码产生的消息更小、更高效可降低磁盘与网络资源占用并提升高吞吐场景性能适用于绝大多数应用json编码产生人类可读的消息便于开发调试也适合与无法处理二进制编码的系统对接。timestamp_mapping.legacy可选trueBoolean是否使用 Avro 中时间戳的旧映射。详见下方「时间戳映射的兼容性」小节。各参数的底层实现说明avro.codec对应源码AVRO_OUTPUT_CODEC配置 key 为codecformat前缀由框架统一加上最终成为avro.codec默认值取DataFileConstants.SNAPPY_CODEC。它只在 Filesystem 的 bulk writer 路径生效AvroFileFormatFactory.java 中的RowDataAvroWriterFactory会调用dataFileWriter.setCodec(CodecFactory.fromString(codec))写入 AvroDataFileWriter的文件头avro.encoding对应源码AVRO_ENCODING枚举AvroEncoding提供BINARY与JSON两种取值。它同时作用于编解码两条路径——AvroSerializationSchema与AvroDeserializationSchema都会读取该配置来选择编码器/解码器timestamp_mapping.legacy对应源码AVRO_TIMESTAMP_LEGACY_MAPPING默认true影响 AvroSchemaConverter.java 中 Flink 时间类型与 Avro 逻辑类型的互转分支。Filesystem 连接器场景下的完整建表示例结合avro.codecCREATE TABLE orders ( order_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3) ) PARTITIONED BY (order_time) WITH ( connector filesystem, path /data/orders, format avro, avro.codec snappy )写入的文件即为标准的 Avro Object Container File文件头携带 schema 与压缩编解码器信息可被 Hive、Spark 等生态直接读取。数据类型映射目前 Avro schema 总是从 table schema 推导尚不支持在 DDL 中显式定义 Avro schema。下表完整列出 Flink SQL 类型到 Avro 类型的映射关系Flink SQL 类型Avro 类型Avro 逻辑类型CHAR / VARCHAR / STRINGstringBOOLEANbooleanBINARY / VARBINARYbytesDECIMALfixeddecimalTINYINTintSMALLINTintINTintBIGINTlongFLOATfloatDOUBLEdoubleDATEintdateTIMEinttime-millisTIMESTAMPlongtimestamp-millisARRAYarrayMAPkey 必须是 string/char/varchar 类型mapMULTISET元素必须是 string/char/varchar 类型mapROWrecord该表与源码 AvroSchemaConverter.java 中convertToSchema(LogicalType, ...)的switch分支一一对应几个值得注意的实现细节DECIMAL在 Flink 侧的底层映射由LogicalTypes.decimal(precision, scale)施加在bytes类型之上表中列为fixed实际编码载体与精度有关转换时会保留 precision/scaleDATE、TIME分别映射为带date、time-millis逻辑类型的int源码对TIME的 precision 做了校验仅支持precision 3即毫秒精度否则抛出IllegalArgumentExceptionMAP与MULTISET统一映射为 Avromap但要求 keyMULTISET则是元素必须是字符串族类型——源码extractValueTypeToAvroMap会显式校验CHARACTER_STRING类型族否则抛出UnsupportedOperationExceptionMULTISET的 value 固定为int计数ROW映射为record嵌套 record 的命名采用rowName_fieldName的递归前缀以保证同一 schema 内 record 名称唯一RAW及其他未列出的类型如INTERVAL等在推导时直接抛出UnsupportedOperationException。Nullable 类型与 union 的映射除了上表列出的类型Flink 支持读写 nullable 的类型。Flink 将 nullable 类型映射为 Avro 的union(something, null)其中something是由 Flink 类型转换得到的 Avro 类型。源码中的nullableSchema(Schema)工具方法正是这一规则的实现private static Schema nullableSchema(Schema schema) { return schema.isNullable() ? schema : Schema.createUnion(SchemaBuilder.builder().nullType(), schema); }反方向同样成立convertToDataType/convertToTypeInfo在处理 AvroUNION时若发现二元 union 中含null则提取出实际类型并标记为 Flink 侧 nullable若是单元 union 则视为非 nullable多于两个分支的 union 无法与 Flink 类型系统直接对齐退化为基于 Kryo 的GENERIC(Object.class)序列化源码注释明确说明。时间戳映射的兼容性timestamp_mapping.legacy这是 Avro format 最容易踩坑的兼容性配置官方文档说明如下在1.19 之前Flink 的默认行为错误地将 SQLTIMESTAMP和TIMESTAMP_LTZ类型都映射到 AvroTIMESTAMP正确的行为是Flink SQLTIMESTAMP应映射 AvroLOCAL TIMESTAMPFlink SQLTIMESTAMP_LTZ应映射 AvroTIMESTAMP。通过将该参数设为false可以禁用旧映射、获得正确映射出于兼容性考虑该参数默认使用旧行为默认值true源码AVRO_TIMESTAMP_LEGACY_MAPPING的defaultValue(true)与此一致。从源码看新旧映射的具体差异AvroSchemaConverter.javalegacy true默认TIMESTAMP无时区→longtimestamp-millis且仅支持precision 3TIMESTAMP_LTZ在推导 schema 时直接抛出UnsupportedOperationException即旧行为下无法正确编码带本地时区的时间戳。legacy false推荐新行为TIMESTAMP→longlocalTimestampMillisprecision ≤ 3或localTimestampMicrosprecision ≤ 6TIMESTAMP_LTZ→longtimestampMillisprecision ≤ 3或timestampMicrosprecision ≤ 6读取方向同样区分Avrotimestamp-millis/micros映射为 FlinkTIMESTAMP_LTZAvrolocalTimestampMillis/micros映射为 FlinkTIMESTAMP。因此如果下游系统如 Confluent Schema Registry 生态遵循标准 Avro 语义timestamp-millis表示带时区的绝对时间点建议显式设置timestamp_mapping.legacy false如果只是与旧版本 Flink 作业读写的数据互通则保持默认即可。源码级工作链路从建表到序列化结合 AvroFormatFactory.java 与 AvroFileFormatFactory.java可以把 Avro format 的运行时链路概括为schema 推导无论是消息格式Kafka 等还是 bulk 文件格式Filesystem都会把表声明的RowType交给AvroSchemaConverter.convertToSchema(...)生成 AvroSchema默认 record 名为org.apache.flink.avro.generated.record消息格式Kafka 等AvroRowDataSerializationSchema.serialize(RowData)内部先RowDataToAvroConverters.createConverter(rowType)将行转成GenericRecord再由嵌套的AvroSerializationSchema按 binary/json 编码输出字节读取方向则对称地由AvroRowDataDeserializationSchema.deserialize(byte[])完成bulk 文件格式Filesystem写入侧RowDataAvroWriterFactory构造DataFileWriterGenericRecord并应用avro.codec压缩读取侧AvroGenericRecordBulkFormat基于文件头内嵌的 schema 读取并与表 schema 做字段级映射源码注释指出文件头 schema 与给定 schema 不一致时reader 会自动按字段名映射参见AbstractAvroBulkFormat的实现。仓库测试 AvroRowDataDeSerializationSchemaTest.java 覆盖了完整的时间戳类型往返如TIMESTAMP(3)字段的序列化/反序列化一致性并显式地将legacyTimestampMapping作为构造参数传入编解码 schema 进行验证可作为阅读上述链路的参考入口。常见注意事项Avro schema 不可显式声明当前版本只能从 table schema 推导。若需要对接已有 Avro schema 的 Topic 或文件请保证表字段名、顺序与 Avro 记录兼容Filesystem 读取时按字段名自动映射消息格式则按声明顺序解析精度限制legacy 映射下TIMESTAMP仅支持毫秒精度precision ≤ 3新映射下TIMESTAMP/TIMESTAMP_LTZ可支持到微秒precision ≤ 6TIME一律只支持毫秒精度超出会抛异常MAP 的 key 限制MAP/MULTISET的 key元素必须是字符串族类型否则 schema 推导直接失败压缩仅对文件系统生效avro.codec是 bulk 文件格式的选项对 Kafka 等消息格式无效消息格式的体积优化应依赖avro.encoding默认 binary以及连接器自身的压缩设置时间戳兼容性跨 1.19 前后版本的数据互通时请显式评估timestamp_mapping.legacy的取值避免时间字段语义错位。以上内容均以当前仓库的 中文文档 与flink-formats/flink-avro模块源码为准。若需要了解 Avro 类型体系的更多细节可查阅官方 Avro 规范中对各类型与逻辑类型的定义。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink JSON Format 完全指南Schema 推导、参数调优与类型映射原理Flink JSON Format 完全指南Schema 推导、参数调优与类型映射原理 JSONJavaScript Object Notation是流式大数据流处理批处理数据工程Apache Flink Table SQL 中的 Parquet 格式依赖、配置参数与数据类型映射详解Apache Flink Table SQL 中的 Parquet 格式依赖、配置参数与数据类型映射详解 导读 本文以 Apache Flink当前仓库 g大数据流处理批处理数据工程Flink CDC StarRocks Pipeline Connector 详解配置参数、Schema 演化机制与类型映射全解Flink CDC StarRocks Pipeline Connector 详解配置参数、Schema 演化机制与类型映射全解 本文以 Flink CDC后端数据集成大数据流处理变更数据捕获数据同步上一篇Obsidian美化终极指南10个CSS代码片段让你的笔记更专业高效下一篇如何在10分钟内完成yuzu模拟器终极部署新手完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考