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-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.1(flink.format.parquet.version),其运行时依赖 Hadoop(hadoop-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指定为filesystem,path指向数据落盘目录(本地路径或 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 | 可选 | false | Boolean | 在 epoch 时间与 LocalDateTime 互转时使用 UTC 时区还是本地时区。Hive 0.x/1.x/2.x 使用本地时区,Hive 3.x 使用 UTC 时区 |
timestamp.time.unit | 可选 | micros | String | 以 int64/LogicalTypes 存储 Parquet 时间戳时的精度单位,取值为nanos/micros/millis |
write.int64.timestamp | 可选 | false | Boolean | 以 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-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:压缩方式,未配置时默认 SNAPPY(CompressionCodecName.SNAPPY.name())。可选值通常包括UNCOMPRESSED、SNAPPY、GZIP、LZO、LZ4、ZSTD、BROTLI等;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 / STRING | BINARY | UTF8 | |
| BOOLEAN | BOOLEAN | ||
| BINARY / VARBINARY | BINARY | ||
| DECIMAL | FIXED_LEN_BYTE_ARRAY | DECIMAL | |
| TINYINT | INT32 | INT_8 | |
| SMALLINT | INT32 | INT_16 | |
| INT | INT32 | ||
| BIGINT | INT64 | ||
| FLOAT | FLOAT | ||
| DOUBLE | DOUBLE | ||
| DATE | INT32 | DATE | |
| TIME | INT32 | TIME_MILLIS | |
| TIMESTAMP | INT96(或 INT64) | ||
| ARRAY | - | LIST | |
| MAP | - | MAP | Parquet 不支持可空 map key |
| MULTISET | - | MAP | Parquet 不支持可空 map key |
| ROW | - | STRUCT |
映射规则的源码级印证
上述映射在 ParquetSchemaConverter.java 中逐类型实现,几个值得深挖的细节:
- 时间戳双模式:当
parquet.write.int64.timestamp为false(默认)时,TIMESTAMP_WITHOUT_TIME_ZONE与TIMESTAMP_WITH_LOCAL_TIME_ZONE统一转换为 int96 原始类型;当为true时,转换为 int64,并按parquet.timestamp.time.unit(nanos/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包下的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/millis); - Hive 版本差异影响时区语义: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 写入 NULL;
- Decimal 精度决定文件字节数:精度越高定长字节越多(
computeMinBytesForDecimalPrecision按2^(8n-1) ≥ 10^p取最小 n),应根据业务实际精度声明字段,避免无谓放大文件体积。
小结
Parquet 格式在 Flink 中承担读写两重角色,核心配置集中在parquet.utc-timezone、parquet.timestamp.time.unit、parquet.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),仅供参考