Flink Parquet 格式全解析:Filesystem 连接器下的读写配置与类型映射实战
Flink Parquet 格式全解析Filesystem 连接器下的读写配置与类型映射实战【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink导读Parquet 是 Apache 大数据生态中最流行的列式存储格式之一在 Flink Table/SQL 中与 Filesystem 连接器配合常被用于数据湖表、数仓 ODS/DWD 层的批式与流式写入。本篇基于 Flink 官方文档 Parquet Format 与仓库内flink-formats/flink-parquet模块源码系统讲解如何在 Flink SQL 中声明 Parquet 表、配置全部格式参数含源码层默认值、理解 Hive/Spark 兼容差异并逐字段对照 Flink 与 Parquet 的类型映射关系。读完本文你将能够独立完成 Parquet 表的建表、读写调优与跨引擎数据交换。Parquet 格式在 Flink 中同时扮演Serialization Schema序列化用于写入与Deserialization Schema反序列化用于读取两种角色即既支持把数据写为.parquet文件也支持把已有 Parquet 文件读回 Flink 表是 Filesystem 连接器 最常用的文件格式之一。依赖引入使用 Parquet 格式需要引入对应的格式依赖。在 Maven 项目中普通 Java 应用DataStream/Table API依赖非 shaded 的flink-parquet模块dependency groupIdorg.apache.flink/groupId artifactIdflink-parquet/artifactId version2.0-SNAPSHOT/version /dependency而 SQL 客户端 / SQL Gateway 等纯 SQL 场景则应使用官方预打包的 shaded 产物flink-sql-parquetjar该 jar 通过 maven-shade-plugin 将flink-parquet、parquet-avro、parquet-hadoop、parquet-format、parquet-column、parquet-encoding、parquet-jackson等依赖一并打入见 flink-sql-parquet/pom.xml放入FLINK_HOME/lib或--jar指定后即可在 SQL 中直接使用format parquet。这一映射关系同样记录在文档站的数据文件 docs/data/sql_connectors.yml 中。从仓库 flink-formats/pom.xml 可以看到当前仓库使用的 Parquet 底层版本为1.13.1flink.format.parquet.version其运行时依赖 Hadoophadoop-common、hadoop-hdfs、hadoop-mapreduce-client-core均以provided作用域提供说明运行环境需自带 Hadoop 依赖。如何创建 Parquet 格式的表下面示例通过 Filesystem 连接器 Parquet 格式创建一张分区表这也是最常见的用法CREATE TABLE user_behavior ( user_id BIGINT, item_id BIGINT, category_id BIGINT, behavior STRING, ts TIMESTAMP(3), dt STRING ) PARTITIONED BY (dt) WITH ( connector filesystem, path /tmp/user_behavior, format parquet )要点说明connector指定为filesystempath指向数据落盘目录本地路径或 HDFS/S3 等分布式文件系统路径均可format指定为parquet格式工厂的factoryIdentifier()正是parquet见 ParquetFileFormatFactory.java该表同时具备读写能力作为 sink 写入时由ParquetRowDataBuilder负责把RowData记录序列化为 Parquet 文件作为 source 读取时由ParquetColumnarRowInputFormat以列式向量化方式读取。写入侧ParquetRowDataBuilder源码会把 SQL 表中声明好的RowType通过ParquetSchemaConverter.convertToParquetMessageType转换为 ParquetMessageType消息结构再据此构建WriteSupport写入每条记录。读取侧则是把 Parquet 的列数据读入 Flink 的列式内存表示RowData/列向量并支持列裁剪projection——createRuntimeDecoder接收投影后的RowType进行按需解码见 ParquetFileFormatFactory.java。Format Options格式参数详解下表为 Parquet 格式的完整参数说明参数是否必填默认值类型说明format必填无String指定使用的格式此处必须为parquetparquet.utc-timezone可选falseBoolean在 epoch 时间与 LocalDateTime 互转时使用 UTC 时区还是本地时区。Hive 0.x/1.x/2.x 使用本地时区Hive 3.x 使用 UTC 时区timestamp.time.unit可选microsString以 int64/LogicalTypes 存储 Parquet 时间戳时的精度单位取值为nanos/micros/milliswrite.int64.timestamp可选falseBoolean以 int64/LogicalTypes 而非 int96/OriginalTypes 写入 Parquet 时间戳。注意此模式下时间戳与时间区无关绝不转换为其他时区注表中parquet.utc-timezone、timestamp.time.unit、write.int64.timestamp这几个带parquet.前缀的键在 SQL 建表语句中书写时同样需要带上parquet.前缀如parquet.utc-timezone true这与下文提到的ParquetOutputFormat参数透传机制保持一致。源码层实现与默认值佐证以上参数的解析集中在ParquetFileFormatFactory源码其内部通过ConfigOption定义了UTC_TIMEZONE键utc-timezoneBoolean 类型默认false与文档表格一致TIMESTAMP_TIME_UNIT键timestamp.time.unit默认microsWRITE_INT64_TIMESTAMP键write.int64.timestamp默认false另有一个文档未列出的BATCH_SIZE键batch-size默认 2048用于控制读取 Parquet 文件时的批大小每批行数可通过parquet.batch-size 4096调整以权衡读取吞吐与内存占用。工厂类在创建读写器时会把所有以parquet.为前缀的格式参数addAllToProperties后批量写入 HadoopConfiguration键为parquet.key再分别传给写入器ParquetRowDataBuilder.createWriterFactory与读取器ParquetColumnarRowInputFormat.createPartitionedFormat。这正是 Parquet 格式支持与 HadoopParquetOutputFormat配置对接的机制基础。透传 ParquetOutputFormat 参数Parquet 格式还支持来自ParquetOutputFormat的配置。由于格式参数最终都会被写入 HadoopConfiguration你可以直接以parquet.*前缀声明任意 Parquet 原生配置例如开启 GZIP 压缩CREATE TABLE user_behavior ( user_id BIGINT, item_id BIGINT, category_id BIGINT, behavior STRING, ts TIMESTAMP(3), dt STRING ) PARTITIONED BY (dt) WITH ( connector filesystem, path /tmp/user_behavior, format parquet, parquet.compression GZIP )从 ParquetRowDataBuilder.java 可以看到创建ParquetWriter时依次读取了以下 Hadoop 配置项ParquetOutputFormat.COMPRESSION压缩方式未配置时默认 SNAPPYCompressionCodecName.SNAPPY.name()。可选值通常包括UNCOMPRESSED、SNAPPY、GZIP、LZO、LZ4、ZSTD、BROTLI等ParquetOutputFormat.BLOCK_SIZErow group 大小ParquetOutputFormat.PAGE_SIZE页大小ParquetOutputFormat.DICTIONARY_PAGE_SIZE字典页大小ParquetOutputFormat.MAX_PADDING_BYTES最大 padding 字节数默认ParquetWriter.MAX_PADDING_SIZE_DEFAULTParquetOutputFormat.ENABLE_DICTIONARY是否启用字典编码ParquetOutputFormat.VALIDATION是否开启写入校验ParquetOutputFormat.WRITER_VERSIONParquet writer 版本。这些参数对控制文件体积、压缩率与下游读取性能有直接影响是 Parquet 表调优时最常触碰的一层配置。数据类型映射Parquet 格式的类型映射当前与 Apache Hive 兼容但默认不与 Apache Spark 兼容主要差异集中在时间戳上Timestamp无论精度如何默认映射为int96Spark 兼容需要通过上文write.int64.timestamp配置项改为写入 int64Decimal按精度映射为定长字节数组FIXED_LEN_BYTE_ARRAY。Flink 类型 → Parquet 类型完整映射表Flink 数据类型Parquet 物理类型Parquet 逻辑类型限制CHAR / VARCHAR / STRINGBINARYUTF8BOOLEANBOOLEANBINARY / VARBINARYBINARYDECIMALFIXED_LEN_BYTE_ARRAYDECIMALTINYINTINT32INT_8SMALLINTINT32INT_16INTINT32BIGINTINT64FLOATFLOATDOUBLEDOUBLEDATEINT32DATETIMEINT32TIME_MILLISTIMESTAMPINT96或 INT64ARRAY-LISTMAP-MAPParquet 不支持可空 map keyMULTISET-MAPParquet 不支持可空 map keyROW-STRUCT映射规则的源码级印证上述映射在 ParquetSchemaConverter.java 中逐类型实现几个值得深挖的细节时间戳双模式当parquet.write.int64.timestamp为false默认时TIMESTAMP_WITHOUT_TIME_ZONE与TIMESTAMP_WITH_LOCAL_TIME_ZONE统一转换为 int96 原始类型当为true时转换为 int64并按parquet.timestamp.time.unitnanos/micros/millis标注LogicalTypeAnnotation.timestampType(false, timeUnit)逻辑类型。写入侧ParquetRowDataWriter同样读取这两个配置来决定时间戳的编码方式源码。需注意int64 模式下的时间戳是时区无关的NEVER converted to a different time zone而 int96 模式配合parquet.utc-timezone决定 epoch 时间与 LocalDateTime 的换算基准——这也是与 Hive 各版本行为差异相关的关键开关Decimal 定长字节数computeMinBytesForDecimalPrecision(precision)从 1 字节起循环计算满足2^(8*bytes-1) 10^precision的最小字节数例如 DECIMAL(10, 2) 需要 5 字节、DECIMAL(18, 2) 需要 8 字节随后以FIXED_LEN_BYTE_ARRAYDECIMAL逻辑类型落盘源码Map/Multiset 的可空 key 处理Parquet 规范不支持可空的 map key因此转换时若 key 类型为可空nullableFlink 会强制copy(false)转为非空类型后再生成 MAP 结构MULTISET 则映射为 key 为元素类型、value 为INT32的 MAP源码ARRAY / ROW分别通过 Parquet 的listOfElements元素统一命名为element与嵌套GroupType生成 LIST / STRUCT 结构。读取侧特性作为 Deserialization SchemaParquet 读取由 ParquetColumnarRowInputFormat.java 与 ParquetSplitReaderUtil.java 等实现具备以下能力列式向量化读取按列批量读取并解码为 Flink 列向量ColumnVector配合batch-size参数控制单批行数显著降低逐行反序列化开销列裁剪projection pushdowncreateRuntimeDecoder接收经过Projection.of(projections)裁剪后的RowType只解码 SQL 查询实际用到的列谓词下推从源码结构看vector/reader包下的BooleanColumnReader、IntColumnReader、LongColumnReader、TimestampColumnReader等实现与 Parquet 页内 RunLength 解码RunLengthDecoder、字典解码ParquetDictionary配合可在页/列块级别跳过无关数据统计信息上报ParquetBulkDecodingFormat实现FileBasedStatisticsReportableInputFormat通过ParquetFormatStatisticsReportUtil.getTableStatistics读取 Parquet 文件页脚中的统计信息为优化器提供TableStats辅助代价估算源码。仓库测试用例 ParquetFileSystemITCase.java 与 ParquetFsStreamingSinkITCase.java 覆盖了文件系统连接器下的端到端读写与流式 Sink 场景可作为理解完整读写链路的最佳入口。最佳实践与注意事项Spark 数据交换前先确认时间戳类型Flink 默认把时间戳写为 int96与 Hive 兼容而 Spark 3 默认按 int64 处理。若要与 Spark 双向读写同一批 Parquet 文件建表时显式设置parquet.write.int64.timestamp true并按需指定parquet.timestamp.time.unit micros或nanos/millisHive 版本差异影响时区语义Hive 0.x/1.x/2.x 使用本地时区解析 epoch 时间Hive 3.x 使用 UTC若跨 Hive 版本读取同一批数据出现时间偏移可通过parquet.utc-timezone true切换转换基准默认false使用本地时区按数据规模选择压缩默认 SNAPPY 在压缩比与 CPU 开销之间较均衡追求更高压缩比可设parquet.compression GZIP或ZSTD需确认运行环境支持对应编解码器追求极致写入吞吐可设UNCOMPRESSED列式存储受益于投影Parquet 天然支持列裁剪与统计信息下推查询时尽量只 SELECT 需要的列并利用分区裁剪如PARTITIONED BY (dt)配合dt 2026-09-22过滤减少扫描量Map/Multiset 键不可空建表时若声明可空的 MAP key 或 MULTISET 元素类型写入端会自动按非空处理业务侧需避免向 key 写入 NULLDecimal 精度决定文件字节数精度越高定长字节越多computeMinBytesForDecimalPrecision按2^(8n-1) ≥ 10^p取最小 n应根据业务实际精度声明字段避免无谓放大文件体积。小结Parquet 格式在 Flink 中承担读写两重角色核心配置集中在parquet.utc-timezone、parquet.timestamp.time.unit、parquet.write.int64.timestamp三个时间戳相关开关与可透传的ParquetOutputFormat参数上类型映射默认对齐 Hiveint96 时间戳 定长字节数组 Decimal需要 Spark 兼容时必须显式开启write.int64.timestamp。结合flink-formats/flink-parquet模块的源码与测试可以进一步按需定制压缩、页大小、批大小等行为将 Parquet 高效地融入 Flink 批流一体的数据湖/数仓实践中。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

搞定导航一下实战项目:3步解决报错堆栈看不懂

搞定导航一下实战项目:3步解决报错堆栈看不懂

搞定导航一下实战项目:3步解决报错堆栈看不懂 昨天帮一个刚入行的兄弟排查代码,他盯着屏幕上的红色报错发呆,说:“大哥,这 StackTrace 一长串英文,我连单词都拼不全,到底哪行代码炸了?”这种场景太常见了。在真实的 实战项目…

2026/9/23 3:13:37 阅读更多 →
3步避开万新神剪手培训坑,一文搞懂公路工程实操

3步避开万新神剪手培训坑,一文搞懂公路工程实操

3步避开万新神剪手培训坑,一文搞懂公路工程实操 官方文档翻了三遍还是云里雾里?别慌,我懂这种抓不住重点的崩溃感。今天不念经,直接带你一文搞懂【万新神剪手】在公路工程微服务里的真实玩法。 概念速懂:它到底在剪什么…

2026/9/23 3:13:37 阅读更多 →
Yii 2 路径别名(Alias)机制完全指南:定义、解析、预定义别名与扩展别名

Yii 2 路径别名(Alias)机制完全指南:定义、解析、预定义别名与扩展别名

Yii 2 路径别名(Alias)机制完全指南:定义、解析、预定义别名与扩展别名 【免费下载链接】yii2 Yii 2: The Fast, Secure and Professional PHP Framework 项目地址: https://gitcode.com/gh_mirrors/yi/yii2 导读 本文以 Yii 2 官方指…

2026/9/23 3:13:37 阅读更多 →

最新新闻

DeepSeek Windows原生部署实战:绕过WSL的高性能方案

DeepSeek Windows原生部署实战:绕过WSL的高性能方案

1. 为什么Windows上部署DeepSeek不是“装个软件”那么简单DeepSeek系列模型(尤其是DeepSeek-V2、DeepSeek-Coder、DeepSeek-MoE等)在开源社区热度持续走高,但很多人点开GitHub仓库看到docker-compose.yml或run.sh脚本时,第一反应是…

2026/9/23 3:54:28 阅读更多 →
Elasticsearch集群变慢?何时该独立部署协调节点及改造方法

Elasticsearch集群变慢?何时该独立部署协调节点及改造方法

说句得罪人的话:大部分人在 Elasticsearch 集群变慢时,第一反应是加数据节点、加副本、加磁盘,很少有人想到“协调节点”这几个字。我见过不少团队,3 个节点扛着每秒几千的查询,CPU 快被打满,业务方天天催&…

2026/9/23 3:54:28 阅读更多 →
GMM与DBSCAN聚类实战对比:突破KMeans瓶颈的概率与密度方法

GMM与DBSCAN聚类实战对比:突破KMeans瓶颈的概率与密度方法

聚类这个问题,平时写代码遇到最多的就是 KMeans,但真正业务里数据一复杂,KMeans 那种"按距离画圆"的思路往往就不够用了。要么簇的形状不规则,要么数据里有明显的离群点,要么样本本身存在重叠,这…

2026/9/23 3:54:28 阅读更多 →
DeepSeek Harness桌面端:智能体工具调用框架与接入实践

DeepSeek Harness桌面端:智能体工具调用框架与接入实践

DeepSeek官方仓库里突然出现了一个叫Harness的桌面端项目,消息在开发者社区传开后,问法五花八门:这跟DeepSeek网页版有什么区别?harness是个框架还是应用?能不能把Codex接进去?为什么还有人把deepseek herm…

2026/9/23 3:54:28 阅读更多 →
12款大模型Three.js代码生成实测:GPT-6 Astra鹈鹕骑车场景夺冠

12款大模型Three.js代码生成实测:GPT-6 Astra鹈鹕骑车场景夺冠

1. 从“鹈鹕骑车”说起:一个被玩坏的经典测试题第一次看到“鹈鹕骑车”这个测试题,大概是在某个深夜刷技术社区的时候。当时的第一反应是:这帮人真会玩。用 Three.js 渲染一只鹈鹕骑自行车的 3D 场景,然后让大模型来生成代码&…

2026/9/23 3:54:28 阅读更多 →
1天重启人生:用24小时重置状态,找回掌控感

1天重启人生:用24小时重置状态,找回掌控感

看到“我悟了!2亿人拜读的万字长文干货,如何在1天内重启你的人生?”这个标题时,我第一反应是:又是一个贩卖焦虑的标题党。毕竟“重启人生”这四个字已经被用滥了,好像只要早起、跑步、列个计划,…

2026/9/23 3:53:28 阅读更多 →

日新闻

3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A…

2026/9/23 0:00:23 阅读更多 →
2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我

2k显示屏性能优化踩坑:版本升级后API全变了,这份源码解析救了我 刚把开发环境的显示器从1080P换到2K,跑老项目直接报错,版本升级后 API…

2026/9/23 0:01:25 阅读更多 →
3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点

3步搞定美眉图实战项目,告别官方文档抓不住重点 官方文档翻了三遍还是云里雾里?别急,美眉图在实战项目中常被用来做数据可视化,但它的原理比你想的简单。今天咱们直接上手,用一个完整的小项目把美眉图跑通,不再死磕那些冗长的理论说明。…

2026/9/23 0:01:25 阅读更多 →

周新闻

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 阅读更多 →