尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
Flink 表格式(Table Formats)全景指南:连接器序列化格式映射与选型实战
Flink 表格式Table Formats全景指南连接器序列化格式映射与选型实战【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink本指南以 Apache Flink Table API / SQL 中的表格式Table Format为绝对主线系统讲解表格式的定义、Flink 内置支持的十余种格式及其与各表连接器的支持矩阵并深入 CSV、JSON 等常用格式的建表实战、参数配置与数据类型映射。读者读完可掌握如何为 Kafka、Filesystem 等连接器正确选择并配置格式理解格式在连接器与运行时之间扮演的二进制数据 ↔ 表列转换角色并能在真实作业中直接套用示例。什么是表格式Table FormatFlink 官方文档对表格式给出了清晰的定义表格式是一种存储格式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 ConnectorsCSVApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、FilesystemJSONApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、Filesystem、ElasticsearchApache AvroApache Kafka、Upsert Kafka、Amazon Kinesis Data Streams、Amazon Kinesis Data Firehose、FilesystemConfluent AvroApache Kafka、Upsert KafkaDebezium CDCApache Kafka、FilesystemCanal CDCApache Kafka、FilesystemMaxwell CDCApache Kafka、FilesystemOGG CDCApache Kafka、FilesystemApache ParquetFilesystemApache ORCFilesystemRawApache 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 ConfigOptionString 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.mdCREATE 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指定要使用的格式这里应为csvcsv.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即抛出错误失败。若忽略字段的解析异常该字段值会被置为nullcsv.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 结果为1E5设为 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.mdCREATE 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声明使用的格式这里应为jsonjson.fail-on-missing-field可选falseBoolean解析字段缺失时是跳过当前字段或行还是抛出错误失败默认 false即抛出错误失败json.ignore-parse-errors可选falseBoolean解析异常时是跳过当前字段或行还是抛出错误失败默认 false。若忽略字段的解析异常该字段值会被置为nulljson.timestamp-format.standard可选SQLString声明输入和输出TIMESTAMP与TIMESTAMP_LTZ的格式支持SQL与ISO-8601SQL以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.123ZISO-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可选FAILString指定处理 Map 中 key 值为空的方法支持FAIL遇到空 key 抛异常、DROP丢弃空 key 数据项、LITERAL用字符串常量替换空 key常量值由json.map-null-key.literal定义json.map-null-key.literal可选nullString当json.map-null-key.mode为LITERAL时指定替换 Map 中空 key 的字符串常量json.encode.decimal-as-plain-number可选falseBoolean将所有 DECIMAL 类型数据保持原状、不使用科学计数法。例0.000000027默认表示为2.7E-8设为 true 时表示为0.000000027json.encode.ignore-null-fields可选falseBoolean仅序列化非 Null 的列默认会序列化所有列无论是否为 Nulldecode.json-parser.enabled可选trueBooleanJsonParser是 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 / STRINGstringBOOLEANbooleanBINARY / VARBINARYstring with encoding: base64DECIMALnumberTINYINTnumberSMALLINTnumberINTnumberBIGINTnumberFLOATnumberDOUBLEnumberDATEstring with format: dateTIMEstring with format: timeTIMESTAMPstring with format: date-timeINTERVALnumberARRAYarrayROWobjectJSON 类型映射Flink SQL 类型JSON 类型CHAR / VARCHAR / STRINGstringBOOLEANbooleanBINARY / VARBINARYstring with encoding: base64DECIMALnumberTINYINTnumberSMALLINTnumberINTnumberBIGINTnumberFLOATnumberDOUBLEnumberDATEstring with format: dateTIMEstring with format: timeTIMESTAMPstring with format: date-timeTIMESTAMP_WITH_LOCAL_TIME_ZONEstring with format: date-time (with UTC time zone)INTERVALnumberARRAYarrayMAP / MULTISETobjectROWobject对比可见两类行式格式对基础类型、日期时间与嵌套结构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),仅供参考
RELATED

相关推荐

LoadRunner性能测试实战:从脚本开发到瓶颈分析全流程

LoadRunner性能测试实战:从脚本开发到瓶颈分析全流程

简介:面向性能测试初学者的一份LoadRunner实验报告,基于Mercury Tours示例应用,适配软件测试课程作业、实验报告撰写及工具自学场景。报告系统地梳理了实验目的与内容,要求掌握脚本录制、编辑与执行技巧,并灵活控制并发…

📅 2026/9/20 15:55:35
Hasura GraphQL Engine 中的托管资源管理:从 `Managed` 到 `ManagedT` 的 Monad 变换器实践

Hasura GraphQL Engine 中的托管资源管理:从 `Managed` 到 `ManagedT` 的 Monad 变换器实践

后端API网关数据库GraphQL 【免费下载链接】graphql-engine Blazing fast, instant realtime GraphQL APIs on all your data with fine grained access control, also trigger webhooks on database events. 项目地址: https://gitcode.com/gh_mirrors/gr/graphql-…

📅 2026/9/20 15:55:35
com.blankj:utilcodex:1.26.0 装不上 Android 12?让走 TaoToken 的 Codex 查

com.blankj:utilcodex:1.26.0 装不上 Android 12?让走 TaoToken 的 Codex 查

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📅 2026/9/20 15:50:33
MORE NEWS

更多资讯

📰

GBrain Ingest 技能深度解析:路由式内容摄取管线与大脑写入契约

GBrain Ingest 技能深度解析:路由式内容摄取管线与大脑写入契约 【免费下载链接】gbrain Garrys Opinionated OpenClaw/Hermes Agent Brain 项目地址: https://gitcode.com/gh_mirrors/gb/gbrain 导读 本文围绕 gbrain 仓库中的 ingest 技能 展开&#xff0…

📰

基于树莓派与ONNX的50Hz双足机器鸭神经控制闭环实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

Luigi 工作流构建指南:深入理解 Task、Target 与 Parameter 三大核心抽象

任务调度工作流自动化批处理后端 【免费下载链接】luigi Luigi is a Python module that helps you build complex pipelines of batch jobs. It handles dependency resolution, workflow management, visualization etc. It also comes with Hadoop support built in. 项目地…

📰

嵌入式AI重构传感器:TinyML在MCU上的工程实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

软考高项论文写作攻略:十大管理领域框架与实战技巧

简介:一份专为软考高级信息系统项目管理师考生设计的十大管理领域论文范文及框架合集,内容紧扣考试要求。资源以XX省公安信息化项目(投资500万元、建设周期1个月)作为贯穿案例,示范如何撰写项目背景、目标、技术选型&a…

📰

Readest 固定版式书籍的 RTL 页面顺序修复:从 5591 看横向从右到左(RTL)阅读的实现与调试

桌面应用跨平台前端 【免费下载链接】readest Readest is a modern, feature-rich ebook reader designed for avid readers offering seamless cross-platform access, powerful tools, and an intuitive interface to elevate your reading experience. 项目地址:…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬