Flink 表格式(Table Formats)全景指南:连接器序列化格式映射与选型实战
2026/9/20 15:55:31 网站建设 项目流程

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)通常同时实现DeserializationFormatFactorySerializationFormatFactory,例如 CsvFormatFactory 正是如此,它同时为运行时提供 CSV 的SerializationSchemaDeserializationSchema实例。

Flink 支持的表格式与连接器支持矩阵

Flink 在表连接器之上提供了一套内置表格式,官方文档以"格式 × 支持的连接器"矩阵的形式给出全景。下表完整收录了当前仓库 overview.md 中列出的格式清单及各自可搭配的连接器:

格式(Format)支持的连接器(Supported Connectors)
CSVApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem
JSONApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem、Elasticsearch
Apache AvroApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem
Confluent AvroApache Kafka、Upsert Kafka
Debezium CDCApache Kafka、Filesystem
Canal CDCApache Kafka、Filesystem
Maxwell CDCApache Kafka、Filesystem
OGG CDCApache Kafka、Filesystem
Apache ParquetFilesystem
Apache ORCFilesystem
RawApache 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。格式工厂随后会:

  1. 通过requiredOptions()/optionalOptions()声明该格式的必选与可选参数(如 JSON 的json.ignore-parse-errors、CSV 的csv.field-delimiter),供FactoryUtil.validateFactoryOptions(...)做校验;
  2. 创建DecodingFormat(读取侧)与EncodingFormat(写入侧)实例;
  3. 由格式实现进一步产出运行时的DeserializationSchema/SerializationSchema(面向消息流式场景)或 bulk 读写接口(面向文件场景)。

值得关注的是 FormatFactory 还提供了forwardOptions()能力:格式可以声明哪些配置项只影响运行时行为(例如时间戳解析格式),可以安全地在作业恢复(plan enrichment)阶段被覆盖,而不会改变执行拓扑。可以看到 JsonFormatFactory 将json.timestamp-format.standardjson.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可选falseBoolean是否禁止对引用的值使用引号(默认 false)。若禁止,则选项'csv.quote-character'不能设置
csv.quote-character可选"String用于围住字段值的引号字符(默认"
csv.allow-comments可选falseBoolean是否允许忽略注释行(默认不允许),注释行以'#'作为起始字符。若允许注释行,请确保csv.ignore-parse-errors也开启从而允许空行
csv.ignore-parse-errors可选falseBoolean解析异常时是跳过当前字段或行,还是抛出错误失败(默认 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可选trueBoolean是否将 BigDecimal 类型数据表示为科学计数法(默认 true)。例如 BigDecimal 值 100000,设为 true 结果为'1E+5',设为 false 结果为100000。注意:仅当值不为 0 且是 10 的倍数时才转为科学计数法

上述参数在源码中对应 CsvFormatFactory 引入的CsvFormatOptions常量(FIELD_DELIMITERALLOW_COMMENTSIGNORE_PARSE_ERRORSNULL_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可选falseBoolean解析字段缺失时,是跳过当前字段或行,还是抛出错误失败(默认 false,即抛出错误失败)
json.ignore-parse-errors可选falseBoolean解析异常时是跳过当前字段或行,还是抛出错误失败(默认 false)。若忽略字段的解析异常,该字段值会被置为null
json.timestamp-format.standard可选'SQL'String声明输入和输出TIMESTAMPTIMESTAMP_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可选falseBoolean将所有 DECIMAL 类型数据保持原状、不使用科学计数法。例:0.000000027默认表示为2.7E-8,设为 true 时表示为0.000000027
json.encode.ignore-null-fields可选falseBoolean仅序列化非 Null 的列,默认会序列化所有列(无论是否为 Null)
decode.json-parser.enabled可选trueBooleanJsonParser是 Jackson 提供的流式读取 JSON 的 API,相比JsonNode方式读取更快、内存消耗更少,且支持嵌套字段的投影下推。默认启用;如遇不兼容问题可禁用并回退到JsonNode方式

从实现上看,JsonFormatFactory 的optionalOptions()与文档参数表一一对应,并且其中json.timestamp-format.standardjson.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 / STRINGstring
BOOLEANboolean
BINARY / VARBINARYstring with encoding: base64
DECIMALnumber
TINYINTnumber
SMALLINTnumber
INTnumber
BIGINTnumber
FLOATnumber
DOUBLEnumber
DATEstring with format: date
TIMEstring with format: time
TIMESTAMPstring with format: date-time
INTERVALnumber
ARRAYarray
ROWobject

JSON 类型映射

Flink SQL 类型JSON 类型
CHAR / VARCHAR / STRINGstring
BOOLEANboolean
BINARY / VARBINARYstring with encoding: base64
DECIMALnumber
TINYINTnumber
SMALLINTnumber
INTnumber
BIGINTnumber
FLOATnumber
DOUBLEnumber
DATEstring with format: date
TIMEstring with format: time
TIMESTAMPstring with format: date-time
TIMESTAMP_WITH_LOCAL_TIME_ZONEstring with format: date-time (with UTC time zone)
INTERVALnumber
ARRAYarray
MAP / MULTISETobject
ROWobject

对比可见:两类行式格式对基础类型、日期时间与嵌套结构(ARRAY/ROW)的映射高度一致,差异主要在于 JSON 额外支持MAP / MULTISETobject的映射,以及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 TABLEWITH子句中声明连接器与格式,即可获得完整的读写能力。

如需进一步深入,可继续阅读本仓库中的下列文档:

  • 格式详情: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),仅供参考

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

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

立即咨询