news 2026/9/23 3:13:30

Flink Parquet 格式全解析:Filesystem 连接器下的读写配置与类型映射实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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> <groupId>org.apache.flink</groupId> <artifactId>flink-parquet</artifactId> <version>2.0-SNAPSHOT</version> </dependency>

而 SQL 客户端 / SQL Gateway 等纯 SQL 场景,则应使用官方预打包的 shaded 产物flink-sql-parquetjar(该 jar 通过 maven-shade-plugin 将flink-parquetparquet-avroparquet-hadoopparquet-formatparquet-columnparquet-encodingparquet-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),其运行时依赖 Hadoop(hadoop-commonhadoop-hdfshadoop-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指定使用的格式,此处必须为'parquet'
parquet.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/millis
write.int64.timestamp可选falseBoolean以 int64/LogicalTypes 而非 int96/OriginalTypes 写入 Parquet 时间戳。注意:此模式下时间戳与时间区无关(绝不转换为其他时区)

:表中parquet.utc-timezonetimestamp.time.unitwrite.int64.timestamp这几个带parquet.前缀的键,在 SQL 建表语句中书写时同样需要带上parquet.前缀(如'parquet.utc-timezone' = 'true'),这与下文提到的ParquetOutputFormat参数透传机制保持一致。

源码层实现与默认值佐证

以上参数的解析集中在ParquetFileFormatFactory(源码),其内部通过ConfigOption定义了:

  • UTC_TIMEZONE:键utc-timezone,Boolean 类型,默认false,与文档表格一致;
  • TIMESTAMP_TIME_UNIT:键timestamp.time.unit,默认"micros"
  • WRITE_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())。可选值通常包括UNCOMPRESSEDSNAPPYGZIPLZOLZ4ZSTDBROTLI等;
  • ParquetOutputFormat.BLOCK_SIZE:row group 大小;
  • ParquetOutputFormat.PAGE_SIZE:页大小;
  • ParquetOutputFormat.DICTIONARY_PAGE_SIZE:字典页大小;
  • ParquetOutputFormat.MAX_PADDING_BYTES:最大 padding 字节数(默认ParquetWriter.MAX_PADDING_SIZE_DEFAULT);
  • ParquetOutputFormat.ENABLE_DICTIONARY:是否启用字典编码;
  • ParquetOutputFormat.VALIDATION:是否开启写入校验;
  • ParquetOutputFormat.WRITER_VERSION:Parquet writer 版本。

这些参数对控制文件体积、压缩率与下游读取性能有直接影响,是 Parquet 表调优时最常触碰的一层配置。

数据类型映射

Parquet 格式的类型映射当前与 Apache Hive 兼容,但默认不与 Apache Spark 兼容,主要差异集中在时间戳上:

  • Timestamp:无论精度如何,默认映射为int96
  • Spark 兼容:需要通过上文write.int64.timestamp配置项改为写入 int64;
  • Decimal:按精度映射为定长字节数组(FIXED_LEN_BYTE_ARRAY)

Flink 类型 → Parquet 类型完整映射表

Flink 数据类型Parquet 物理类型Parquet 逻辑类型限制
CHAR / VARCHAR / STRINGBINARYUTF8
BOOLEANBOOLEAN
BINARY / VARBINARYBINARY
DECIMALFIXED_LEN_BYTE_ARRAYDECIMAL
TINYINTINT32INT_8
SMALLINTINT32INT_16
INTINT32
BIGINTINT64
FLOATFLOAT
DOUBLEDOUBLE
DATEINT32DATE
TIMEINT32TIME_MILLIS
TIMESTAMPINT96(或 INT64)
ARRAY-LIST
MAP-MAPParquet 不支持可空 map key
MULTISET-MAPParquet 不支持可空 map key
ROW-STRUCT

映射规则的源码级印证

上述映射在 ParquetSchemaConverter.java 中逐类型实现,几个值得深挖的细节:

  • 时间戳双模式:当parquet.write.int64.timestampfalse(默认)时,TIMESTAMP_WITHOUT_TIME_ZONETIMESTAMP_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_ARRAY+DECIMAL逻辑类型落盘(源码);
  • Map/Multiset 的可空 key 处理:Parquet 规范不支持可空的 map key,因此转换时若 key 类型为可空(nullable),Flink 会强制copy(false)转为非空类型后再生成 MAP 结构,MULTISET 则映射为 key 为元素类型、value 为INT32的 MAP(源码);
  • ARRAY / ROW:分别通过 Parquet 的listOfElements(元素统一命名为element)与嵌套GroupType生成 LIST / STRUCT 结构。

读取侧特性

作为 Deserialization Schema,Parquet 读取由 ParquetColumnarRowInputFormat.java 与 ParquetSplitReaderUtil.java 等实现,具备以下能力:

  • 列式向量化读取:按列批量读取并解码为 Flink 列向量(ColumnVector),配合batch-size参数控制单批行数,显著降低逐行反序列化开销;
  • 列裁剪(projection pushdown)createRuntimeDecoder接收经过Projection.of(projections)裁剪后的RowType,只解码 SQL 查询实际用到的列;
  • 谓词下推(从源码结构看)vector/reader包下的BooleanColumnReaderIntColumnReaderLongColumnReaderTimestampColumnReader等实现与 Parquet 页内 RunLength 解码(RunLengthDecoder)、字典解码(ParquetDictionary)配合,可在页/列块级别跳过无关数据;
  • 统计信息上报ParquetBulkDecodingFormat实现FileBasedStatisticsReportableInputFormat,通过ParquetFormatStatisticsReportUtil.getTableStatistics读取 Parquet 文件页脚中的统计信息,为优化器提供TableStats辅助代价估算(源码)。

仓库测试用例 ParquetFileSystemITCase.java 与 ParquetFsStreamingSinkITCase.java 覆盖了文件系统连接器下的端到端读写与流式 Sink 场景,可作为理解完整读写链路的最佳入口。

最佳实践与注意事项

  1. Spark 数据交换前先确认时间戳类型:Flink 默认把时间戳写为 int96(与 Hive 兼容),而 Spark 3 默认按 int64 处理。若要与 Spark 双向读写同一批 Parquet 文件,建表时显式设置'parquet.write.int64.timestamp' = 'true'并按需指定'parquet.timestamp.time.unit' = 'micros'(或nanos/millis);
  2. Hive 版本差异影响时区语义:Hive 0.x/1.x/2.x 使用本地时区解析 epoch 时间,Hive 3.x 使用 UTC;若跨 Hive 版本读取同一批数据出现时间偏移,可通过'parquet.utc-timezone' = 'true'切换转换基准(默认false使用本地时区);
  3. 按数据规模选择压缩:默认 SNAPPY 在压缩比与 CPU 开销之间较均衡;追求更高压缩比可设'parquet.compression' = 'GZIP''ZSTD'(需确认运行环境支持对应编解码器),追求极致写入吞吐可设'UNCOMPRESSED'
  4. 列式存储受益于投影:Parquet 天然支持列裁剪与统计信息下推,查询时尽量只 SELECT 需要的列,并利用分区裁剪(如PARTITIONED BY (dt)配合dt = '2026-09-22'过滤)减少扫描量;
  5. Map/Multiset 键不可空:建表时若声明可空的 MAP key 或 MULTISET 元素类型,写入端会自动按非空处理,业务侧需避免向 key 写入 NULL;
  6. Decimal 精度决定文件字节数:精度越高定长字节越多(computeMinBytesForDecimalPrecision2^(8n-1) ≥ 10^p取最小 n),应根据业务实际精度声明字段,避免无谓放大文件体积。

小结

Parquet 格式在 Flink 中承担读写两重角色,核心配置集中在parquet.utc-timezoneparquet.timestamp.time.unitparquet.write.int64.timestamp三个时间戳相关开关与可透传的ParquetOutputFormat参数上;类型映射默认对齐 Hive(int96 时间戳 + 定长字节数组 Decimal),需要 Spark 兼容时必须显式开启write.int64.timestamp。结合flink-formats/flink-parquet模块的源码与测试,可以进一步按需定制压缩、页大小、批大小等行为,将 Parquet 高效地融入 Flink 批流一体的数据湖/数仓实践中。

【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 3:13:12

搭建4k影院项目:3步搞定视频流性能优化

搭建4k影院项目:3步搞定视频流性能优化 面试被问原理答不上来?别慌。很多开发者能写出业务代码,但一提到4k影院这种高负载场景的性能优化,就卡壳了。这不仅是技术深度的试金石,更是你从“码农”进阶为“架构师”的必经之路。…

作者头像 李华
网站建设 2026/9/23 3:12:49

告别有妖气下载报错:手写完整示例破解技术难题

告别有妖气下载报错:手写完整示例破解技术难题 看了一堆教程还是不会写项目?别急着骂自己笨,大概率是你没看懂底层逻辑。很多开发者卡在“有妖气下载”这类资源获取脚本上,不是语法不会,而是没搞懂请求拦截、数据解析和文件落盘的完整闭环。今天不整虚的,直接给你一份能跑通的 完整示例…

作者头像 李华
网站建设 2026/9/23 3:12:11

热力图工具选型:行为还原精度与多端埋点实战指南

1. 热力图不是“看热闹”&#xff0c;而是用户行为的X光片你点开一个热力图工具&#xff0c;看到页面上红红绿绿的色块&#xff0c;第一反应可能是&#xff1a;“哇&#xff0c;这块好热&#xff01;”——但真正用过三年以上、带过五个以上产品团队的从业者会立刻问三个问题&a…

作者头像 李华
网站建设 2026/9/23 3:11:55

STM32串口通信面试避坑指南附完整示例

STM32串口通信面试避坑指南附完整示例 版本升级后 API 全变了?别慌,这往往是嵌入式开发者从“会调库”到“懂底层”的分水岭。很多候选人拿着 HAL 库的代码去面试,被问到底层寄存器配置就卡壳,或者在 HAL 库迁移到 LL 库时一脸茫然。今天这篇 stm32串口通信…

作者头像 李华
网站建设 2026/9/23 3:11:50

qqtang源码拆解:3个核心技巧教你彻底搞懂底层逻辑保姆级教程

qqtang源码拆解:3个核心技巧教你彻底搞懂底层逻辑保姆级教程 看了一堆教程还是不会写项目?别慌,今天这篇 保姆级教程 带你直接扒开源码看。 很多开发者陷入一个怪圈:文档看了三遍,代码抄了两遍,一到实战就懵。问题出在哪?你只学了“怎么调”,没搞懂“为什么这么调”。以 qqtang…

作者头像 李华
网站建设 2026/9/23 3:11:41

盛名来电通源码解析:搞懂3个高频考点,面试不再慌

盛名来电通源码解析:搞懂3个高频考点,面试不再慌 版本升级后 API 全变了,文档还跟不上,你盯着屏幕发呆时,面试官正盯着你的简历问:“说说盛名来电通底层是怎么处理并发请求的?”别慌。今天咱们不背八股文,直接拆解【盛名来电通】的【源码解析】,用真实代码带你通关。 考点梳理:面试官到底想考什么…

作者头像 李华