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),仅供参考

相关新闻

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

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

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

2026/9/22 21:59:40 阅读更多 →
PDF题库结构化:Python实现ERP考试知识图谱构建

PDF题库结构化:Python实现ERP考试知识图谱构建

简介:本资源是面向金蝶云星空认证考生与ERP实施顾问的备考核心资料,聚焦K/3 Cloud系统高频考点与实操难点,覆盖采购管理、许可申请、核算设置、工作流配置、SQL库管理、套打设置、报销单处理、科目初始化、基础资料分配、销售预收、用户同步、…

2026/9/22 12:59:38 阅读更多 →
计算机联锁仿真系统设计全解析:站场建模、联锁引擎与故障注入

计算机联锁仿真系统设计全解析:站场建模、联锁引擎与故障注入

简介:计算机联锁仿真系统软件设计文档,面向铁路信号、自动化及相关专业学生与工程技术人员,主要讲解以VC开发古浪车站上行咽喉联锁仿真系统的原理与实现。文档从计算机联锁系统的基本结构和功能入手,详细阐述进路建立、进路解锁等…

2026/9/20 15:54:32 阅读更多 →

最新新闻

全金属机甲斗神怎么打:配置环境卡半天后的最佳实践

全金属机甲斗神怎么打:配置环境卡半天后的最佳实践

全金属机甲斗神怎么打:配置环境卡半天后的最佳实践 配置环境就卡半天,这是很多开发者在接触新框架或复杂系统时的第一道坎。面对全金属机甲斗神怎么打这个看似与编程无关的问题,实则隐喻了我们在处理高复杂度、多依赖、强耦合系统时的痛点。很多教程只讲理…

2026/9/22 21:59:21 阅读更多 →
DNF天帷禁地通关全解:完整示例拆解底层逻辑

DNF天帷禁地通关全解:完整示例拆解底层逻辑

DNF天帷禁地通关全解:完整示例拆解底层逻辑 官方文档里关于副本机制的说明往往晦涩难懂,几十页的文本让人抓不住重点。别慌,我们直接切入核心,用一套 完整示例…

2026/9/22 21:59:20 阅读更多 →
2026最新头像文字源码解析:面试被问原理答不上来?

2026最新头像文字源码解析:面试被问原理答不上来?

2026最新头像文字源码解析:面试被问原理答不上来? 面试被问到“头像文字”底层渲染逻辑,答不上来?这不仅是技术盲区,更是2026最新前端工程化能力的试金石。很多开发者停留在 avatar…

2026/9/22 21:59:20 阅读更多 →
属于c高频面试题

属于c高频面试题

3个实战项目带你彻底搞懂C语言指针属于谁 版本升级后 API 全变了,这是很多老程序员的噩梦,也是新手入门时的第一道坎。 别慌,今天不聊虚的。我们直接上手一个【实战项目】,通过解决一个真实的内存管理问题,来彻底搞懂那个让人头秃的问题:…

2026/9/22 21:59:20 阅读更多 →
3个坑教你手写实现装饰设计培训项目

3个坑教你手写实现装饰设计培训项目

3个坑教你手写实现装饰设计培训项目 版本升级后 API 全变了,昨天还能跑的装饰工程数据接口,今天全报 404。别急着骂娘,这其实是底层逻辑变了。很多从业者还在死记硬背旧版参数,结果被新版校验机制卡得死死的。与其天天查文档改参数,不如直接手…

2026/9/22 21:58:20 阅读更多 →
qq播放器下载源码拆解:3个实战项目级技巧

qq播放器下载源码拆解:3个实战项目级技巧

qq播放器下载源码拆解:3个实战项目级技巧 学会语法却不知怎么搭项目,是大多数开发者转行或进阶时的最大卡点。很多人背下了 Python 的类继承、Java 的并发包,甚至刷完了 LeetCode 的前 200 题,但面对一个真实的…

2026/9/22 21:58:20 阅读更多 →

日新闻

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天 配置环境就卡半天?别怪机器慢,多半是你没选对工具链。在Java、Go或Python的项目现场, 手写实现…

2026/9/22 0:00:41 阅读更多 →
剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑 面试被问原理答不上来,是不是常态?别慌。很多开发者对着 GitHub 开源仓库里的代码发呆,看似简单实则暗藏玄机。今天这份【剑帝加点】速查手册,直接带你拆解核心实现,把面试必考的原理讲透。…

2026/9/22 0:00:41 阅读更多 →
手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优 复制来的代码跑不通不知道怎么调?别慌,这种“复制粘贴地狱”在开发圈太常见了。尤其是做 图片压缩网站…

2026/9/22 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/22 4:32:41 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/22 4:38:57 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/22 8:51:04 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/21 15:36:51 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/21 15:36:51 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/22 2:43:42 阅读更多 →