ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

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),仅供参考
返回列表