Flink 表格式(Table Formats)全景指南:连接器序列化格式映射与选型实战
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
本指南以 Apache Flink Table API / SQL 中的"表格式(Table Format)"为绝对主线,系统讲解表格式的定义、Flink 内置支持的十余种格式及其与各表连接器的支持矩阵,并深入 CSV、JSON 等常用格式的建表实战、参数配置与数据类型映射。读者读完可掌握如何为 Kafka、Filesystem 等连接器正确选择并配置格式,理解格式在连接器与运行时之间扮演的"二进制数据 ↔ 表列"转换角色,并能在真实作业中直接套用示例。
什么是表格式(Table Format)
Flink 官方文档对表格式给出了清晰的定义:表格式是一种存储格式(storage format),它定义了如何把二进制数据(binary data)映射到表的列(table columns)上。它与"连接器(Connector)"是正交的两个概念——连接器负责接入外部系统(如 Kafka、文件系统),格式则负责解释或生成这些系统里流动的字节流。
从源码角度可以进一步印证这一抽象:在 Format.java 中,格式被描述为"连接器格式的基接口",并且可以从两个维度进行区分:
- 应用上下文:格式作用于
DynamicTableSource(读取侧)还是DynamicTableSink(写入侧); - 运行时实现接口:格式最终需要产出哪种运行时实现,例如
DeserializationSchema(反序列化)或某种 bulk 接口。
对应地,源码中将格式细分为 DecodingFormat(把外部二进制数据解码为RowData,供 Source 读取)与 EncodingFormat(把RowData编码为外部二进制数据,供 Sink 写出)两类能力。一个格式工厂(Format Factory)通常同时实现DeserializationFormatFactory与SerializationFormatFactory,例如 CsvFormatFactory 正是如此,它同时为运行时提供 CSV 的SerializationSchema和DeserializationSchema实例。
Flink 支持的表格式与连接器支持矩阵
Flink 在表连接器之上提供了一套内置表格式,官方文档以"格式 × 支持的连接器"矩阵的形式给出全景。下表完整收录了当前仓库 overview.md 中列出的格式清单及各自可搭配的连接器:
| 格式(Format) | 支持的连接器(Supported Connectors) |
|---|---|
| CSV | Apache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem |
| JSON | Apache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem、Elasticsearch |
| Apache Avro | Apache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem |
| Confluent Avro | Apache Kafka、Upsert Kafka |
| Debezium CDC | Apache Kafka、Filesystem |
| Canal CDC | Apache Kafka、Filesystem |
| Maxwell CDC | Apache Kafka、Filesystem |
| OGG CDC | Apache Kafka、Filesystem |
| Apache Parquet | Filesystem |
| Apache ORC | Filesystem |
| Raw | Apache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem |
从矩阵中可以提炼出几条关键规律:
- Filesystem 连接器的格式支持最全(本仓库中对应文档为 filesystem.md),从面向行的 CSV/JSON,到面向列的 Parquet/ORC,再到各类 CDC 格式均可使用,这与文件系统"按文件存储、格式自解释"的特性一致;
- Apache Kafka / Upsert Kafka 是格式覆盖最广的消息类连接器,几乎支持上表全部格式;
- Elasticsearch 连接器仅与 JSON 格式搭配(文档中未列出其他格式);
- Parquet 与 ORC 只服务于 Filesystem,因为它们本质上是列式文件存储格式,天然面向批量文件场景;
- 此外,当前仓库的格式目录中还提供了 Protobuf 的独立文档页,其实现位于 flink-protobuf 模块,并配套有 flink-sql-protobuf 的 SQL 打包模块。
格式如何被连接器发现与装配:Factory 机制
在 Flink Table 体系中,WITH子句里的'format' = 'xxx'是连接器与格式之间的"装配开关"。该选项在源码 FactoryUtil.java 中被定义为:
public static final ConfigOption<String> FORMAT = ConfigOptions.key("format") ...FactoryUtil 会按format的值(或key.format/value.format这类带后缀的变体)去发现对应的格式工厂(FormatFactory)。每个格式工厂都通过factoryIdentifier()声明自己的标识符,例如 JsonFormatFactory 中:
public static final String IDENTIFIER = "json";也就是说,SQL 中写'format' = 'json'时,正是通过该标识符匹配到JsonFormatFactory。格式工厂随后会:
- 通过
requiredOptions()/optionalOptions()声明该格式的必选与可选参数(如 JSON 的json.ignore-parse-errors、CSV 的csv.field-delimiter),供FactoryUtil.validateFactoryOptions(...)做校验; - 创建
DecodingFormat(读取侧)与EncodingFormat(写入侧)实例; - 由格式实现进一步产出运行时的
DeserializationSchema/SerializationSchema(面向消息流式场景)或 bulk 读写接口(面向文件场景)。
值得关注的是 FormatFactory 还提供了forwardOptions()能力:格式可以声明哪些配置项只影响运行时行为(例如时间戳解析格式),可以安全地在作业恢复(plan enrichment)阶段被覆盖,而不会改变执行拓扑。可以看到 JsonFormatFactory 将json.timestamp-format.standard、json.map-null-key.mode等解析相关参数声明为 forward 选项——修改这些参数不会影响 ChangelogMode 等拓扑级能力。
实战一:CSV 格式 + Kafka 连接器建表
CSV 格式允许基于 CSV schema 解析和生成 CSV 数据,当前 CSV schema 由 table schema 推断而来,不支持显式定义 CSV schema。以下建表示例完整引自 csv.md:
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' = 'csv', 'csv.ignore-parse-errors' = 'true', 'csv.allow-comments' = 'true' )CSV 格式参数一览
| 参数 | 是否必选 | 默认值 | 类型 | 描述 |
|---|---|---|---|---|
format | 必选 | (none) | String | 指定要使用的格式,这里应为'csv' |
csv.field-delimiter | 可选 | , | String | 字段分隔符(默认','),必须为单字符。可使用反斜杠指定特殊字符,如'\t'代表制表符;也可通过 unicode 编码在纯 SQL 文本中指定,如'csv.field-delimiter' = U&'\0001'代表0x01字符 |
csv.disable-quote-character | 可选 | false | Boolean | 是否禁止对引用的值使用引号(默认 false)。若禁止,则选项'csv.quote-character'不能设置 |
csv.quote-character | 可选 | " | String | 用于围住字段值的引号字符(默认") |
csv.allow-comments | 可选 | false | Boolean | 是否允许忽略注释行(默认不允许),注释行以'#'作为起始字符。若允许注释行,请确保csv.ignore-parse-errors也开启从而允许空行 |
csv.ignore-parse-errors | 可选 | false | Boolean | 解析异常时是跳过当前字段或行,还是抛出错误失败(默认 false,即抛出错误失败)。若忽略字段的解析异常,该字段值会被置为null |
csv.array-element-delimiter | 可选 | ; | String | 分隔数组和行元素的字符串(默认';') |
csv.escape-character | 可选 | (none) | String | 转义字符(默认关闭) |
csv.null-literal | 可选 | (none) | String | 指定识别为 null 值的字符串(默认禁用)。输入端将该字符串转为 null 值,输出端将 null 值转成该字符串 |
csv.write-bigdecimal-in-scientific-notation | 可选 | true | Boolean | 是否将 BigDecimal 类型数据表示为科学计数法(默认 true)。例如 BigDecimal 值 100000,设为 true 结果为'1E+5',设为 false 结果为100000。注意:仅当值不为 0 且是 10 的倍数时才转为科学计数法 |
上述参数在源码中对应 CsvFormatFactory 引入的CsvFormatOptions常量(FIELD_DELIMITER、ALLOW_COMMENTS、IGNORE_PARSE_ERRORS、NULL_LITERAL等),建表时设置的每个csv.*键都会被逐一映射到这些配置项并参与校验。
实战二:JSON 格式 + Kafka 连接器建表
JSON 格式能读写 JSON 格式的数据,当前 JSON schema 同样从 table schema 自动推导,不支持显式定义。以下建表示例完整引自 json.md:
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' = 'json', 'json.fail-on-missing-field' = 'false', 'json.ignore-parse-errors' = 'true' )JSON 格式参数一览
| 参数 | 是否必选 | 默认值 | 类型 | 描述 |
|---|---|---|---|---|
format | 必选 | (none) | String | 声明使用的格式,这里应为'json' |
json.fail-on-missing-field | 可选 | false | Boolean | 解析字段缺失时,是跳过当前字段或行,还是抛出错误失败(默认 false,即抛出错误失败) |
json.ignore-parse-errors | 可选 | false | Boolean | 解析异常时是跳过当前字段或行,还是抛出错误失败(默认 false)。若忽略字段的解析异常,该字段值会被置为null |
json.timestamp-format.standard | 可选 | 'SQL' | String | 声明输入和输出TIMESTAMP与TIMESTAMP_LTZ的格式,支持'SQL'与'ISO-8601':'SQL'以yyyy-MM-dd HH:mm:ss.s{precision}解析 TIMESTAMP(如2020-12-30 12:13:14.123),以yyyy-MM-dd HH:mm:ss.s{precision}'Z'解析 TIMESTAMP_LTZ(如2020-12-30 12:13:14.123Z);'ISO-8601'以yyyy-MM-ddTHH:mm:ss.s{precision}解析 TIMESTAMP(如2020-12-30T12:13:14.123),以yyyy-MM-ddTHH:mm:ss.s{precision}'Z'解析 TIMESTAMP_LTZ;输出均与输入格式保持一致 |
json.map-null-key.mode | 可选 | 'FAIL' | String | 指定处理 Map 中 key 值为空的方法,支持'FAIL'(遇到空 key 抛异常)、'DROP'(丢弃空 key 数据项)、'LITERAL'(用字符串常量替换空 key,常量值由'json.map-null-key.literal'定义) |
json.map-null-key.literal | 可选 | 'null' | String | 当'json.map-null-key.mode'为LITERAL时,指定替换 Map 中空 key 的字符串常量 |
json.encode.decimal-as-plain-number | 可选 | false | Boolean | 将所有 DECIMAL 类型数据保持原状、不使用科学计数法。例:0.000000027默认表示为2.7E-8,设为 true 时表示为0.000000027 |
json.encode.ignore-null-fields | 可选 | false | Boolean | 仅序列化非 Null 的列,默认会序列化所有列(无论是否为 Null) |
decode.json-parser.enabled | 可选 | true | Boolean | JsonParser是 Jackson 提供的流式读取 JSON 的 API,相比JsonNode方式读取更快、内存消耗更少,且支持嵌套字段的投影下推。默认启用;如遇不兼容问题可禁用并回退到JsonNode方式 |
从实现上看,JsonFormatFactory 的optionalOptions()与文档参数表一一对应,并且其中json.timestamp-format.standard、json.map-null-key.*、json.encode.*等被声明为forwardOptions(),说明它们属于"只影响运行时解析行为、不影响拓扑"的稳定选项,可以被安全地覆盖。
数据类型映射:Flink 类型与外部格式类型的对应关系
CSV 与 JSON 格式均基于 table schema 自动推导 schema,其序列化/反序列化在底层使用 jackson databind API 解析与生成数据。两个格式的类型映射表如下(分别完整引自 csv.md 与 json.md)。
CSV 类型映射
| Flink SQL 类型 | CSV 类型 |
|---|---|
CHAR / VARCHAR / STRING | string |
BOOLEAN | boolean |
BINARY / VARBINARY | string with encoding: base64 |
DECIMAL | number |
TINYINT | number |
SMALLINT | number |
INT | number |
BIGINT | number |
FLOAT | number |
DOUBLE | number |
DATE | string with format: date |
TIME | string with format: time |
TIMESTAMP | string with format: date-time |
INTERVAL | number |
ARRAY | array |
ROW | object |
JSON 类型映射
| Flink SQL 类型 | JSON 类型 |
|---|---|
CHAR / VARCHAR / STRING | string |
BOOLEAN | boolean |
BINARY / VARBINARY | string with encoding: base64 |
DECIMAL | number |
TINYINT | number |
SMALLINT | number |
INT | number |
BIGINT | number |
FLOAT | number |
DOUBLE | number |
DATE | string with format: date |
TIME | string with format: time |
TIMESTAMP | string with format: date-time |
TIMESTAMP_WITH_LOCAL_TIME_ZONE | string with format: date-time (with UTC time zone) |
INTERVAL | number |
ARRAY | array |
MAP / MULTISET | object |
ROW | object |
对比可见:两类行式格式对基础类型、日期时间与嵌套结构(ARRAY/ROW)的映射高度一致,差异主要在于 JSON 额外支持MAP / MULTISET到object的映射,以及TIMESTAMP_LTZ的 UTC 时区语义;而BINARY / VARBINARY在两种格式中都以 base64 字符串承载。在设计表结构时,应确保外部数据(CSV 文件、JSON 消息)的实际形态与上表一致,避免隐式类型不匹配导致的解析失败。
格式选型建议
结合上文的支持矩阵与各格式特点,可以按以下维度进行选型:
- 流式消息场景(Kafka 等):首选 CSV / JSON / Avro。CSV 与 JSON 对 schema 要求宽松、可直接由 table schema 推导,适合快速接入;Avro 适合需要强 schema 管理、与上游 Hadoop/流生态深度集成的场景;Confluent Avro 则适用于使用 Confluent Schema Registry 管理 schema 的 Kafka 生态;
- 数据库变更捕获(CDC)场景:根据上游 CDC 工具选择对应格式——Debezium CDC、Canal CDC、Maxwell CDC、OGG CDC,它们均以 JSON 为载体描述行级变更(insert/update/delete),并支持 Kafka 与 Filesystem 两类连接器;
- 批量文件 / 数仓场景(Filesystem):面向列的 Apache Parquet 与 Apache ORC 是首选,具备高压缩比与列裁剪优势;需要保留原始字节时可用 Raw 格式;
- 简单二进制透传:Raw 格式适合单列、无结构解析的裸字节场景,同样覆盖 Kafka、Upsert Kafka、Kinesis、Firehose、Filesystem 等主流连接器。
小结与延伸阅读
表格式是 Flink Table 生态中连接"外部存储的二进制形态"与"表列的逻辑结构"的关键抽象:连接器负责传输与落盘,格式负责映射与解析,二者通过'format'选项和 Format Factory 机制在运行时完成装配。开发者只需在CREATE TABLE的WITH子句中声明连接器与格式,即可获得完整的读写能力。
如需进一步深入,可继续阅读本仓库中的下列文档:
- 格式详情:CSV、JSON、Apache Avro、Confluent Avro、Protobuf、Debezium CDC、Canal CDC、Maxwell CDC、OGG CDC、Apache Parquet、Apache ORC、Raw;
- 连接器详情:Filesystem;
- 源码参考:Format.java、DecodingFormat.java、EncodingFormat.java、FormatFactory.java、CsvFormatFactory、JsonFormatFactory。
【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考