Parquet 诞生时,Hadoop 生态中的数据处理框架正在快速增多。同一份数据会被不同工具反复读取,却缺少一种独立于计算框架的高效列式格式。分析查询通常只需要少数几列;按完整记录读取,会把大量 I/O 浪费在无关字段上。
Parquet 并非第一个列式存储格式。它的价值在于把列式表示、编码和压缩写进一套公开的文件规范,并借鉴 Google Dremel 的记录拆分与重组算法来表示嵌套数据。项目最初由 Twitter、Cloudera 等贡献者开发,2014 年进入 Apache Incubator,次年毕业成为 Apache 顶级项目。现在,多种语言和分析工具都能直接读写 Parquet。
正文固定读取 parquet-lab/transactions-small.parquet。这是一份未加密的普通 Parquet 文件,由附录 B中的脚本生成。样例有 30,000 行、3 个 Row Group 和 8 个叶子列。精确大小与校验值集中列在附录中。除附录 A 的嵌套示例外,正文中的 parquet 与 DuckDB 命令都读取这份文件。
现在给查询引擎一个具体问题:从这份交易文件中,按用户汇总 2026 年 5 月 8 日 03:00 到 04:00 的交易净额。
SELECT user_id, SUM(amount)
FROM transactions
WHERE created_at >= TIMESTAMP '2026-05-08 03:00:00'
AND created_at < TIMESTAMP '2026-05-08 04:00:00'
GROUP BY user_id;通过 SQL 可得到结果列、过滤条件和聚合方式,却无法告诉读取器这些列数据在文件中的具体位置。执行查询前还要回答三个问题:
- Schema 如何定义,
user_id、amount和created_at对应哪些叶子列? - 哪些 Row Group 和 Page 不可能满足时间条件,可以直接跳过?
- 候选列采用什么编码和压缩,读取器怎样找到并解码它们?
1. Parquet 文件结构#
先看完整文件,再模拟查询。图中的粗边界表示 File -> Row Group -> Column Chunk -> Page。细线只表示局部放大,不是数据流。
这几层各有分工。Row Group 划分一批行,Column Chunk 保存这批行中的一个叶子列,Page 再保存该列的一小段编码数据。每层的大小都由写入器选择,不是格式常量。
Bloom Filter 和 Page Index 不在这条包含链中。Bloom Filter 用来快速判断“某个值一定不存在”。Page Index 由 ColumnIndex 和 OffsetIndex 组成:前者像页级范围目录,记录每页的边界;后者像位置表,记录页所在的字节范围和行区间。
这个样例写入器的物理布局如下;它描述的是本文件的写入顺序,而不是规范强制的可选辅助结构排列顺序。
PAR1
Row Group 0 ... Row Group N
optional Bloom Filter / ColumnIndex / OffsetIndex byte ranges
serialized FileMetaData
4-byte FileMetaData length
PAR1本文样例以 PAR1 开头和结尾。数据区保存一个或多个 Row Group,每个 Row Group 又为每个 Schema 叶子列保存一个 Column Chunk。Column Chunk 是连续的字节区,里面可以先写 Dictionary Page,再写多个 Data Page。Page 是局部编码、压缩、解压和校验的边界。
这里的布局只适用于未加密文件。Parquet Modular Encryption 的 plaintext-footer 模式仍使用 PAR1。encrypted-footer 模式在文件首尾使用 PARE,并在加密 Footer 前增加 FileCryptoMetaData。后文读取末尾 8 字节并验证 PAR1,说的都是本文样例。
叶子列是嵌套结构展开后的存储路径。例如,
a: {b: {c: 1}}会展开成a.b.c,这样底层的值才能按列保存。
Bloom Filter、ColumnIndex 和 OffsetIndex 都服务于单个 Column Chunk,但字节不在该 Chunk 的连续 Page 区域内。读取器要先解析文件尾的 FileMetaData,再根据其中的引用找到它们。因此,这张图描述文件层级,不描述读取顺序。
2. 用 SQL 驱动一次读取#
查询引擎解析开篇 SQL 后,会把投影和过滤条件交给 Parquet 扫描算子。投影需要 user_id 和 amount,过滤需要 created_at,因此扫描阶段要读三个叶子列。
GROUP BY user_id 和 SUM(amount) 由上层执行算子完成。Parquet 读取器只负责返回通过时间条件的列数据,不负责聚合。
先区分格式能力与工具的实际行为。样例生成器会写出 Page Index,parquet-cli 也能解析它。不过,PyArrow 25.0.0 文档明确说明,其读取端尚不使用 Page Index。
DuckDB 1.5.5 的 Parquet 扫描器会用 Row Group Statistics 和 Bloom Filter 做裁剪。它没有读取 column_index_offset 或 offset_index_offset。因此,下表第 5 阶段只是格式层模拟,不是本次 DuckDB 查询的执行跟踪。
读取过程可以按下面的阶段理解:
| 阶段 | 开篇 SQL 带来的要求 | 读取器的动作 | 本样例的结果 |
|---|---|---|---|
| 1 | 打开 transactions 文件 |
读取末尾 8 字节,定位并解析 FileMetaData | 得到 Schema、3 个 Row Group 及其物理引用 |
| 2 | WHERE created_at >= 03:00 AND created_at < 04:00 |
比较 created_at 的 Column Chunk Statistics |
排除 RG0 和 RG2,只保留 RG1 |
| 3 | 读取候选记录批次 | 进入 RG1,并按 Schema 解析其中的叶子列 | RG1 对应 10,000 条候选记录 |
| 4 | 输出 user_id,聚合 amount,过滤 created_at |
只选择这三个 Column Chunk | 其余五列不进入本次扫描 |
| 5 | 继续执行时间过滤 | 若读取器支持 Page Index,用 ColumnIndex 筛选 Page,再按 OffsetIndex 定位字节和行区间 |
格式元数据表明 Page 0 至 Page 4 仍可能命中;这不是 DuckDB 实测路径 |
| 6 | GROUP BY user_id、SUM(amount) |
解码候选数据,执行原始时间谓词,再把投影列交给聚合算子 | DuckDB 实测有 3,597 行进入聚合,得到 180 个用户 |
后面的章节沿这条路径展开:FileMetaData -> 谓词下推 -> Row Group -> Column Chunk -> Page。Statistics、Bloom Filter 和 Page Index 会在各自参与裁剪的位置解释。
3. FileMetaData:从文件尾开始#
读取这份未加密文件时,不必从文件头一路扫描到目标列。读取器先取最后 8 字节:前 4 字节记录 FileMetaData 的长度,后 4 字节是末尾魔数 PAR1。文件大小减去这 8 字节,再减去元数据长度,就是 FileMetaData 的起点。
encrypted-footer 模式使用 PARE 和另一套 Footer 解析流程,不适用下面的步骤。
3.1 定位并解析 FileMetaData#
长度字段使用小端序。本文样例解析出的 FileMetaData 长度是 4,186 字节,起点是 1,484,874。原始尾部字节、换算公式和命令集中列在附录 B。
FileMetaData 使用 Thrift Compact Protocol 序列化。parquet-cli 将读取并解析 FileMetaData 的子命令命名为 footer:
parquet footer --raw parquet-lab/transactions-small.parquet解析后的结构可以概括为:
FileMetaData
├── schema: SchemaElement[]
└── row_groups: RowGroup[]
└── columns: ColumnChunk[]
├── meta_data: ColumnMetaData
│ ├── type / encodings / codec / num_values / sizes
│ ├── data_page_offset / dictionary_page_offset
│ ├── statistics: Statistics
│ └── bloom_filter_offset / bloom_filter_length
├── column_index_offset / column_index_length
└── offset_index_offset / offset_index_lengthFileMetaData 不直接装入所有辅助结构。它主要保存 Schema、Row Group、Column Chunk,以及通往其他字节区域的引用。读取器解析这些引用后,才会跳到具体的 Page 或索引。
引用分属两个层次:
| 所属结构 | 保存的内容 |
|---|---|
ColumnChunk |
ColumnIndex 与 OffsetIndex 的 offset/length |
ColumnMetaData |
Dictionary/Data Page 的 offset、Chunk 大小、Statistics 和 Bloom Filter 引用 |
file_offset 是已废弃字段,历史实现对其含义并不一致。读取器不能把样例中的 file_offset = 0 当作 Column Chunk 起点,而应使用内层的 Page offset 与 Chunk 大小。RG1 / created_at 的相关引用值见附录 B。
到这里,读取器只拿到了 Schema、Statistics 和物理引用,还没有读取任何 Data Page。
3.2 用 Schema 解析 SQL 中的列#
parquet schema 从刚刚解析的 FileMetaData 中读取 schema: SchemaElement[],将结果打印成便于阅读的 message 结构:
parquet schema --parquet parquet-lab/transactions-small.parquetmessage schema {
required int64 id;
required int32 mch_id;
required int64 user_id;
required binary out_trans_no (STRING);
required int64 amount;
required int64 balance;
optional int64 created_at (TIMESTAMP(MICROS,false));
optional int64 updated_at (TIMESTAMP(MICROS,false));
}开篇 SQL 投影列为 user_id 和 amount,使用 created_at 谓词过滤,因此扫描算子只需继续追踪这三条叶子路径即可。
Schema 还规定字段是否可空,以及它的物理表示和逻辑语义。created_at 的 Physical Type 是 INT64,Logical Type 是 TIMESTAMP(MICROS,false)。读取器因此把整数解释为微秒精度、无 UTC 调整语义的时间戳。out_trans_no 的 Physical Type 是 BYTE_ARRAY,其 STRING 注解表示这些字节采用 UTF-8。
扫描算子现在已经确定 created_at、user_id 和 amount 三条叶子路径。下一步先看 created_at 的 Statistics 能否排除整批记录。
4. 谓词下推:Statistics 与 Bloom Filter#
查询引擎会把开篇 SQL 中的 created_at 条件下推给 Parquet 扫描算子。“下推”不是让文件执行 SQL,而是让读取器在解码前先看元数据,排除肯定不可能命中的 Row Group。剩余候选仍要读取,并执行原始谓词。
4.1 Statistics 的字段结构#
Statistics 可以理解成一段列数据的摘要。读取器先看摘要中的边界和空值数量,再决定是否值得读取数据本身。它是一个所有字段都可选的 Thrift 结构,主要字段如下:
| 字段 | 含义 |
|---|---|
min_value / max_value |
首选的最小、最大边界;新写入器使用这组字段 |
null_count |
null 数量;字段缺失不等于 0 |
distinct_count |
不同值数量,主要用于基数估算 |
is_min_value_exact / is_max_value_exact |
对应边界是否为数据中实际出现的极值 |
nan_count |
浮点列中的 NaN 数量;字段缺失时要假定可能存在 NaN |
min / max |
已废弃的旧边界字段 |
边界值使用 PLAIN 编码。读取器还要结合物理类型、逻辑类型和 ColumnOrder 解释它们,不能把所有原始字节按同一种方式排序。完整 Thrift 定义见 Parquet Format 2.13.0。
Statistics 有两种常见粒度。ColumnMetaData.statistics 描述整个 Column Chunk,也就是一个 Row Group 中的某个叶子列。Data Page Header 中的 Statistics 只描述单页。
开篇查询先用前者裁剪 Row Group。第 7 章的 ColumnIndex 是另一套逐页结构,不是 Statistics 的别名。
4.2 用 Statistics 裁剪 Row Group#
created_at 的 Statistics 位于每个 Row Group 对应的 ColumnMetaData 中。parquet meta 会把物理时间戳解码成人能直接比较的 min/max。完整输出还包含另外七列,下面只摘录 Row Group 标题和 created_at:
parquet meta parquet-lab/transactions-small.parquetRow group 0: count: 10000 47.53 B records start: 4 total(compressed): 464.189 kB total(uncompressed):941.346 kB
...
created_at INT64 S _ R 10000 7.65 B 10 "2026-05-08T00:00:00.000000" / "2026-05-08T02:46:39.000000"
Row group 1: count: 10000 47.52 B records start: 475334 total(compressed): 464.092 kB total(uncompressed):946.834 kB
...
created_at INT64 S _ R 10000 7.65 B 10 "2026-05-08T02:46:40.000000" / "2026-05-08T05:33:19.000000"
Row group 2: count: 10000 47.47 B records start: 950564 total(compressed): 463.551 kB total(uncompressed):951.110 kB
...
created_at INT64 S _ R 10000 7.65 B 10 "2026-05-08T05:33:20.000000" / "2026-05-08T08:19:59.000000"S _ R 不是 Parquet 规范里的字段,而是 parquet-cli 的紧凑显示:
| 符号 | 在这段输出中的含义 |
|---|---|
S |
Column Chunk 使用 SNAPPY |
_ |
Dictionary Page 使用 PLAIN |
R |
Data Page 使用 RLE_DICTIONARY |
下划线在其他输出位置也可能表示“未压缩”,所以不能脱离列头单独解读。第 7.4 节还会看到这种情况。
将这些范围与 SQL 的半开区间 [03:00:00, 04:00:00) 比较:
Row Group 0: 00:00:00 .. 02:46:39 -> skip
Row Group 1: 02:46:40 .. 05:33:19 -> candidate
Row Group 2: 05:33:20 .. 08:19:59 -> skip
Predicate: 03:00:00 .. 04:00:00RG0 的最大值早于查询下界,RG2 的最小值不小于查询上界,因此两组都不可能命中。RG1 的范围与谓词相交,只能标记为候选。读取器仍要检查其中的记录。nulls = 10 不会改变结论,因为 SQL 的范围比较不会让 null 通过过滤。
Statistics 只能用于排除,min/max 相交并不证明数据一定命中。字段缺失,或者边界经过截断且不再可信时,读取器必须回退到读取数据,再执行剩余过滤。
4.3 Bloom Filter:另一条排除路径#
Bloom Filter 是一种概率型成员测试结构。它能确定某个值不存在,却只能说某个值“可能存在”,因此适合等值过滤,不适合这里的时间范围条件。
开篇 SQL 没有按 out_trans_no 做等值过滤,所以 Bloom Filter 不参与这次执行。样例仍为该列写入了 Bloom Filter,下面用它观察这两种返回结果:
parquet bloom-filter \
-c out_trans_no \
-v '102117463212-017782057860228-1-0-281-1778205786-0,missing-transaction-000' \
parquet-lab/transactions-small.parquetRow group 0:
value 102117463212-017782057860228-1-0-281-1778205786-0 maybe exists.
value missing-transaction-000 NOT exists.
Row group 1:
value 102117463212-017782057860228-1-0-281-1778205786-0 NOT exists.
value missing-transaction-000 NOT exists.
Row group 2:
value 102117463212-017782057860228-1-0-281-1778205786-0 NOT exists.
value missing-transaction-000 NOT exists.maybe exists 不能证明值存在,NOT exists 才允许读取器排除对应 Row Group。范围查询仍依赖 Statistics 或 Page Index,不使用 Bloom Filter。
5. Row Group#
Statistics 已经把候选范围缩小到 RG1。Row Group 是一批完整记录的物理边界。组内包含同一段行号范围的全部叶子列,不同 Row Group 可以独立读取。
parquet meta 报告 RG1 有 10,000 条记录,压缩后的全部 Column Chunk 合计约 464.092 kB:
parquet meta parquet-lab/transactions-small.parquetRow group 1: count: 10000 47.52 B records start: 475334 total(compressed): 464.092 kB total(uncompressed):946.834 kB样例生成器按顺序写入三个 RecordBatch,因此 RG1 对应逻辑行 10,000..19,999。这个范围来自生成顺序,Parquet 并没有额外保存一列全局行号。start: 475334 也是 parquet-cli 针对本文件报告的位置,不是规范要求的固定布局。
RG1 仍然包含八个叶子列。查询不需要把它们全部读入内存,下一步会按 SQL 的投影和谓词选出三个 Column Chunk。
6. Column Chunk:列裁剪#
RG1 按 Schema 拆成八个 Column Chunk,每个 Chunk 连续保存一个叶子列的 Page。开篇 SQL 只需要 user_id、amount 和 created_at。其余五列可以跳过。
同一条 parquet meta 命令会列出 RG1 的八个 Column Chunk。下面只保留本次扫描涉及的三行:
parquet meta parquet-lab/transactions-small.parquetRow group 1: count: 10000 47.52 B records start: 475334 total(compressed): 464.092 kB total(uncompressed):946.834 kB
--------------------------------------------------------------------------------
type encodings count avg size nulls min / max
user_id INT64 S _ R 10000 0.36 B 0 "102117463712" / "102117464211"
amount INT64 S _ R 10000 0.12 B 0 "-250000" / "1250000"
created_at INT64 S _ R 10000 7.65 B 10 "2026-05-08T02:46:40.000000" / "2026-05-08T05:33:19.000000"这里的 S _ R 与上一节含义相同:SNAPPY、PLAIN Dictionary Page 和 RLE_DICTIONARY Data Page。三个 Chunk 的压缩后大小差异很大,但读取器不必靠顺序扫描猜测边界。它会分别使用 dictPageOffset、dataPageOffset 和 totalSize 定位。本文样例的具体值集中列在附录 B。
7. Page:编码与解压边界#
Page 是 Column Chunk 内的编码、压缩、解压和校验边界,但它不等于一次 I/O 请求。顺序扫描时,读取器常会合并读取多个 Page,甚至预取更大的连续区间。只有实现支持 Page Index 等定位信息时,Range Read 才可能收窄到少数 Page。
还要区分 Page 的外层结构和页内类型。所有 Page 都采用 PageHeader + Page body,再由 PageHeader.type 决定页专属 Header 和页体布局。Page Index 是用于裁剪 Data Page 的独立元数据,不是另一层 Page 容器。
7.1 Page 序列与通用 PageHeader#
一个 Column Chunk 在物理上是一串连续的 Page。使用字典编码时,开头最多有一张 Dictionary Page,后面是一张或多张 Data Page。没有字典时,Column Chunk 可以直接从 Data Page 开始。Page 数量和目标大小由写入器决定,不是格式常量。
每张 Page 都由两段连续字节组成。前面是 Thrift 序列化的 PageHeader,后面是长度为 compressed_page_size 的 Page body。PageHeader 不是定长结构,读取器必须先完成反序列化,才能知道它占了多少字节。
没有
headerSize,读取器怎么知道PageHeader在哪里结束?Thrift Compact Protocol 会逐字段解析
PageHeader,并在当前 struct 的T_STOP处结束。它不是在原始字节中搜索某个分隔符。字段的 wire type 还允许读取器跳过不认识的可选字段。解析结束时,输入游标正好指向 Page body。工具显示的
headerSize是“当前游标减去 Page 起点”的结果,不是文件里另存的字段。逐步验算见附录 B。
PageHeader 保存 type、页体压缩前后的大小、可选 CRC,以及与 type 对应的页专属 Header。这里的大小不包含 PageHeader 自身。页专属 Header 是嵌套字段,不是紧跟在通用 Header 后面的另一段物理 Header。
PageType |
PageHeader 中设置的字段 |
Page body 的含义 |
|---|---|---|
DATA_PAGE = 0 |
data_page_header |
Data Page V1 的 RL、DL 和编码值 |
INDEX_PAGE = 1 |
index_page_header |
格式枚举中的 Index Page;本文样例未使用 |
DICTIONARY_PAGE = 2 |
dictionary_page_header |
编码后的字典值 |
DATA_PAGE_V2 = 3 |
data_page_header_v2 |
分段保存的 RL、DL 和编码值 |
INDEX_PAGE 不能与下一节的 Page Index 混为一谈。用于裁剪的 Page Index 指 ColumnIndex 和 OffsetIndex 两个独立结构,它们的位置和长度由 ColumnChunk 元数据引用。
CRC 字段位于 PageHeader 内,校验范围却只覆盖写入磁盘的 Page body,不包含 Header。读取器先反序列化通用 Header,再根据 type 解释页专属 Header 和后续 compressed_page_size 字节。
7.2 Page Index:ColumnIndex 与 OffsetIndex#
格式层已经把范围缩小到 RG1 的三个 Column Chunk。支持 Page Index 的读取器还能继续裁剪;不支持的读取器会读取整个候选 Chunk,再执行原始谓词。
Page Index 是服务于单个 Column Chunk 的可选逐页元数据。两个组成部分分工如下:
| 结构 | 保存什么 | 解决的问题 |
|---|---|---|
ColumnIndex |
各 Data Page 的 min/max、null 和边界顺序 | 哪些 Page 仍可能命中 |
OffsetIndex |
各 Data Page 的偏移、总长度和首行索引 | Page 在哪里、覆盖哪些行 |
ColumnIndex 若存在,对应的 OffsetIndex 必须存在;OffsetIndex 也可以单独出现。两者都只描述 Data Page,不包含开头可能存在的 Dictionary Page。
查询区间是 [03:00:00, 04:00:00)。样例的 Page 0 至 Page 4 与它相交,仍是候选。Page 5 从 04:10:00 开始,后面的页都可以排除。这里展示的是格式能提供的信息;DuckDB 1.5.5 并没有执行这一步。
ColumnIndex 中与边界相邻的几页可以简化成下面这张表:
| Data Page | min | max | 与查询区间的关系 |
|---|---|---|---|
| Page 0 | 02:46:40 | 03:03:19 | 相交,保留 |
| Page 4 | 03:53:20 | 04:09:59 | 相交,保留 |
| Page 5 | 04:10:00 | 04:26:39 | 不相交,排除 |
读取器先用 ColumnIndex 找候选页,再用 OffsetIndex 把页号换成字节范围和 RG 内行区间。它还可以据此寻找 user_id、amount 中覆盖相同行区间的 Page。完整的 offset、压缩长度和首行索引见附录 B。
Page Index 缺失,或者读取器没有实现这条路径,都不会影响正确性。代价只是读取更大的候选范围,并在解码后执行原始谓词。
7.3 读取 Dictionary Page 与 Data Page#
沿 RG1 / created_at 的 Column Chunk 往里看,先是一张 Dictionary Page,随后是十张 Data Page。pages --raw 会解析共用的 PageHeader,再按 type 展开对应的页专属 Header。
| 物理顺序 | type |
页专属 Header | 作用 |
|---|---|---|---|
| Page 0 | 2 |
dictionary_page_header |
保存按 PLAIN 编码的字典值 |
| Page 1 | 3 |
data_page_header_v2 |
保存 levels 和字典 ID |
两张 Page 的 Header 序列化后长度不同,说明 PageHeader 不是固定宽度前缀。Header 大小、Page 起点和页体长度见附录 B。
created_at 有 10,000 个逻辑位置,其中 10 个是 null,其余 9,990 个时间戳互不相同。这是高基数数据。样例统一设置了 use_dictionary=True,因此仍然写出 Dictionary Page,方便观察字典 ID 与 Page offset。这不代表生产环境也该这样配置。写入器可以关闭字典,也可以在字典达到限制后让后续 Data Page 回退到 PLAIN。
Data Page V2 的页体依次保存 RL、DL 和 encoded values。前两段的长度直接写在 DataPageHeaderV2 中,并保持未压缩;is_compressed 只控制 values 区。即使 values 没有压缩,compressed_page_size 这个字段名也不会改变,它仍表示磁盘上的整个 Page body 长度。
CRC 用于发现意外损坏,不提供加密、签名或真实性证明。CRC 值写在 Header 中,校验范围则是 Page body。样例中的分段字节数和 CRC 值也放在附录表格中。
本文样例通过 PyArrow 的 data_page_version="2.0" 选择 V2。V1 与 V2 共用 PageHeader + Page body 外层结构,主要差别如下:
| 对比项 | Data Page V1 | Data Page V2 |
|---|---|---|
| levels | RL、DL 与 values 连续写入页体 | Header 直接记录 RL、DL 的字节长度 |
| Codec 范围 | 启用后整体压缩 | RL/DL 不压缩,只有 values 受 is_compressed 控制 |
| 提前读 levels | 通常要先解压整个页体 | 可以不解压 values,先读取 levels |
V2 没有取代 V1。两者是 PageType 中并存的 Data Page 类型,V2 也并非在所有场景都更合适。
7.4 编码与压缩#
编码改变值的表示,例如把重复值换成较短的字典 ID。Codec 再压缩编码后的字节。两步可以独立选择:SNAPPY 不是 RLE_DICTIONARY,使用 Dictionary 编码也不代表页体一定经过压缩。
| 编码 | 保存什么 | 适用场景 |
|---|---|---|
| PLAIN | 按 Physical Type 直接表示值;Dictionary Page 也用它保存字典值 | 通用回退、字典本体 |
| RLE / Bit Packing | 连续相同值或窄位宽整数 | Definition/Repetition Level、小整数 ID |
| Dictionary + RLE_DICTIONARY | Dictionary Page 存去重值,Data Page 存字典 ID | 低基数字符串、枚举式数据 |
mch_id 是最直观的例子。样例的 30,000 行都属于商户 102,因此每个 Row Group 的 Dictionary Page 只有一个值。后面的十张 Data Page 各保存 1,000 个字典 ID。
parquet pages -c mch_id parquet-lab/transactions-small.parquetpage type enc count size rows
0-D dict S _ 1 4 B
0-1 data _ R 1000 4 B 1000
0-2 data _ R 1000 4 B 1000
...
0-10 data _ R 1000 4 B 1000这段 Page 输出把 Codec 状态和值编码并排缩写。Dictionary Page 的 S _ 表示 SNAPPY + PLAIN。Data Page V2 的 _ R 表示 values 区未压缩,并使用 RLE_DICTIONARY。这里的下划线含义取决于所在列,不能单独解读。
还有几种常见但不必逐项展开的值编码:
| 编码 | 典型数据形状 |
|---|---|
DELTA_BINARY_PACKED |
相邻差值较小的整数、递增序列 |
DELTA_BYTE_ARRAY |
相邻值共享较长前缀的字节串 |
BYTE_STREAM_SPLIT |
浮点数或定长字节值,便于后续压缩 |
这些编码属于格式能力,本样例没有实际使用它们。样例能够直接验证的是 Dictionary Page、RLE_DICTIONARY、Data Page V2 及其 Page 数。
本文样例通过 max_rows_per_page=1000 限制每页的最大行数,这不是 Parquet 格式的固定值。CRC 也是可选功能;未开启时,读取器仍须正确读取文件。
8. 完整 Parquet 文件与查询结果#
前面依次查看了 FileMetaData、Statistics、Row Group、Column Chunk 和 Page。现在把物理层级与元数据引用放回同一张图。左侧是文件包含关系,右侧是 FileMetaData 字段树;中间的虚线表示引用,不是数据流。
样例文件的精确偏移图和对应数值放在附录 B,避免在主线中反复穿插字节验算。
开篇 SQL 对这些结构的使用顺序可以压缩成一条完整链路:
- 读取末尾 8 字节,根据长度字段定位 FileMetaData。
- 从 Schema 解析
created_at、user_id和amount,其余五个叶子列不进入扫描。 - 用
created_atStatistics 排除 RG0、RG2,保留 RG1。 - 沿 RG1 的三个 Column Chunk 引用跳到各自的 Page 区域。
- 若读取器实现 Page Index 裁剪,可用
ColumnIndex保留created_at的 Data Page 0 至 Page 4。随后用OffsetIndex取得字节范围和行区间。本文的 DuckDB 1.5.5 查询没有走这条路径。 - 读取器校验、解压、解码候选页,执行剩余时间过滤,再把
user_id与amount交给聚合算子。
物理元数据只能说明哪些数据可以读取或跳过。最后还要用独立实现核对查询结果。下面的 CTE 保留开篇 SQL 的过滤和聚合,同时把结果压成一行,便于复现时比较:
duckdb -csv -c "
WITH filtered AS (
SELECT user_id, amount
FROM read_parquet('parquet-lab/transactions-small.parquet')
WHERE created_at >= TIMESTAMP '2026-05-08 03:00:00'
AND created_at < TIMESTAMP '2026-05-08 04:00:00'
), result AS (
SELECT user_id, SUM(amount) AS net_amount
FROM filtered
GROUP BY user_id
)
SELECT (SELECT COUNT(*) FROM filtered) AS matched_rows,
COUNT(*) AS users,
SUM(net_amount) AS total_net_amount,
MIN(user_id) AS min_user_id,
MAX(user_id) AS max_user_id
FROM result;
"matched_rows,users,total_net_amount,min_user_id,max_user_id
3597,180,449165000,102117463752,102117463931时间条件命中 3,597 行,聚合后得到 180 个用户。DuckDB 用来核对逻辑结果,parquet-cli 用来检查 Schema、Statistics、偏移和 Page Header。两类证据解决的问题不同。
9. 总结#
Parquet 的读取过程可以归纳为以下几点:
- Parquet 是面向分析型负载的列式文件格式。它定义文件结构,不负责事务、查询执行或存储管理。
- 单个文件按 File、Row Group、Column Chunk 和 Page 分层。Row Group 划分记录批次,Column Chunk 保存一个叶子列,Page 是局部编码、压缩和校验的边界。
- FileMetaData 写在数据之后,使未加密 Parquet 可以单遍顺序写入。文件完成并写出尾部长度与魔数后,读取器才能定位 Footer。
- Statistics 与 Bloom Filter 先排除不可能命中的范围。Page Index 再用
ColumnIndex选择页,用OffsetIndex定位字节和行区间。这些结构都是优化,缺失或未被读取器使用都不能改变查询结果。 - 编码和压缩是两个步骤。Dictionary、RLE 与 Delta 改变值的表示,SNAPPY、ZSTD 等 Codec 再压缩编码后的字节。
附录 A:嵌套数据与 DL/RL#
正文查询只涉及扁平列。遇到 struct 或 list 时,仅保存叶子值还不够,读取器还要知道 null 出现在哪一层,以及多个值是否属于同一条记录。Parquet 用 Definition Level(DL)和 Repetition Level(RL)保存这些边界。
嵌套样例只有四条逻辑记录,但每一种“空”都不同:
{"event_id": 1, "user": {"city": "Shanghai"}, "actions": [{"type": "bet", "amount": -100}, {"type": "win", "amount": 250}]}
{"event_id": 2, "user": null, "actions": []}
{"event_id": 3, "user": {"city": null}, "actions": null}
{"event_id": 4, "user": {"city": "Beijing"}, "actions": [{"type": "bet", "amount": null}]}这四条记录覆盖了几种容易混淆的状态:user=null 是 null struct,city=null 是 struct 存在但 leaf 为空。actions=[] 是 empty list,actions=null 则是 null list。最后一条记录还有一个存在但 amount=null 的元素。只把值排成数组,会丢失这些差别。
先看 parquet-cli 从 FileMetaData 还原的三层 LIST Schema:
parquet schema --parquet parquet-lab/nested-events.parquetmessage schema {
required int64 event_id;
optional group user {
optional binary city (STRING);
}
optional group actions (LIST) {
repeated group list {
required group element {
required binary type (STRING);
optional int64 amount;
}
}
}
}最终写入的是四条叶子路径,而不是一个不可拆分的 actions blob:
event_id
user.city
actions.list.element.type
actions.list.element.amountDL 表示从根走到叶子时,可选或重复节点定义到了多深。RL 表示当前条目与前一个条目共享到哪一层重复路径。
null、empty list 等情况仍会产生 level 条目,但不会向物理 values 区写入占位值。按上面四条记录展开,得到:
| 叶子列 | max DL / RL | DL 序列 | RL 序列 | 编码后的非 null 值流 |
|---|---|---|---|---|
user.city |
2 / 0 |
2, 0, 1, 2 |
0, 0, 0, 0 |
Shanghai, Beijing |
actions.list.element.type |
2 / 1 |
2, 2, 1, 0, 2 |
0, 1, 0, 0, 0 |
bet, win, bet |
actions.list.element.amount |
3 / 1 |
3, 3, 1, 0, 2 |
0, 1, 0, 0, 0 |
-100, 250 |
第二条记录的 list 存在但没有元素,所以停在 DL 1。第三条连 list 都不存在,只到 DL 0。第四条的 list 和 element 都存在,只缺少可选的 amount,因此到 DL 2,但不会向物理 values 区写入整数。
Data Page 的 num_values 统计 level 条目,其中包含 null 和空列表标记。它不等于编码后的非 null 值数量。
RL 只在第一条记录的第二个 action 变为 1,表示它仍处在同一条记录的 actions 重复层内。其余位置都从一条新记录开始。
Page V2 确实保存了两段 level 字节流:
parquet pages --raw \
-c actions.list.element.amount \
parquet-lab/nested-events.parquetPage 1. (offset: 374, headerSize: 28)
{
"crc" : -587734138,
"data_page_header_v2" : {
"definition_levels_byte_length" : 3,
"num_rows" : 2,
"num_values" : 3,
"repetition_levels_byte_length" : 2
},
"type" : 3
}这个输出只能证明 DL/RL 区域存在,不会打印解码后的 level 数组。表中的序列由固定 Schema 和四条记录逐项推导,再通过 DuckDB 的跨实现读取结果核对:
duckdb -box -c "
SELECT *
FROM read_parquet('parquet-lab/nested-events.parquet')
ORDER BY event_id;
"event_id user actions
1 {'city': Shanghai} [{'type': bet, 'amount': -100}, {'type': win, 'amount': 250}]
2 NULL []
3 {'city': NULL} NULL
4 {'city': Beijing} [{'type': bet, 'amount': NULL}]parquet-cli 1.16.0 的 cat / scan 通过 Avro 投影读取这份三层 LIST 时,会在 GroupColumnIO.getLast 报错。Schema 与 Page 检查不受影响。因此,本文用 parquet-cli 检查物理结构,用 DuckDB 1.5.5 验证逻辑记录,不用 parquet cat 验证这份嵌套样例。
附录 B:可复现实验与检查清单#
样例字节与偏移明细#
正文只保留结构与查询主线。本节集中列出样例文件的偏移、长度和校验值,供复现或排错时查阅。
文件尾的定位数据如下:
| 项目 | 值 | 计算或含义 |
|---|---|---|
| 文件大小 | 1,489,068 字节 |
wc -c 的结果 |
| 末尾 8 字节起点 | 1,489,060 |
1,489,068 - 8 |
| 末尾 8 字节 | 5a 10 00 00 50 41 52 31 |
前 4 字节是长度,后 4 字节是 PAR1 |
| FileMetaData 长度 | 4,186 字节 |
小端序 5a 10 00 00,即 0x105a |
| FileMetaData 起点 | 1,484,874 |
1,489,068 - 8 - 4,186 |
RG1 中三个查询列的 Column Chunk 引用如下:
| 列 | Dictionary Page 起点 | 第一张 Data Page 起点 | Chunk 压缩后总长度 |
|---|---|---|---|
user_id |
561,790 | 563,829 | 3,619 |
amount |
725,493 | 725,595 | 1,182 |
created_at |
797,552 | 858,004 | 76,506 |
RG1 / created_at 还引用了两个 Page Index 结构。它的 ColumnIndex 位于 1,479,014,长度为 244 字节;OffsetIndex 位于 1,483,672,长度为 113 字节。作为对照,RG0 / out_trans_no 的 Bloom Filter 位于 1,425,240,长度为 16,401 字节。
Dictionary Page 和第一张 Data Page 可以连续验算:
| 对象 | 起点 | Header | Page body | 合计或下一位置 |
|---|---|---|---|---|
created_at Dictionary Page |
797,552 |
26 字节 | 60,426 字节 | 下一页:797,552 + 26 + 60,426 = 858,004 |
created_at Data Page 0 |
858,004 |
32 字节 | 1,261 字节 | OffsetIndex 总长度:32 + 1,261 = 1,293 |
PageHeader 的读取步骤也可以按游标位置理解:
| 步骤 | 读取器动作 | 结果 |
|---|---|---|
| 1 | 从 Page 起点按 Thrift Compact Protocol 逐字段解析 | 游标跨过变长 Header |
| 2 | 在当前 struct 的 T_STOP 处结束 |
游标指向 Page body 首字节 |
| 3 | 用“当前游标 - Page 起点”计算 headerSize |
得到 26 或 32 字节 |
| 4 | 再读取 compressed_page_size 字节 |
得到完整 Page body |
Dictionary Page 含 9,990 个 INT64 字典值。PLAIN 编码前的大小是 9,990 * 8 = 79,920 字节,SNAPPY 将 Page body 压缩为 60,426 字节。
第一张 Data Page V2 的 Page body 分段如下:
| Page body 范围 | 内容 | 大小与状态 |
|---|---|---|
[0, 0) |
Repetition Levels | 0 字节;该列没有重复路径 |
[0, 8) |
Definition Levels | 8 字节;V2 中不压缩 |
[8, 1261) |
RLE_DICTIONARY 编码的字典 ID | 1,253 字节;is_compressed = 0,不压缩 |
三段合计 1,261 字节,所以 compressed_page_size 与 uncompressed_page_size 都是 1,261。该页的 crc = 1278766728,校验范围是这 1,261 字节 Page body,不包含 Header。
Page Index 中与查询边界相邻的位置如下:
| Data Page | offset | 压缩后总长度 | RG 内首行 | min | max |
|---|---|---|---|---|---|
| Page 0 | 858,004 |
1,293 | 0 | 02:46:40 | 03:03:19 |
| Page 4 | 863,800 |
1,668 | 4,000 | 03:53:20 | 04:09:59 |
| Page 5 | 865,468 |
1,668 | 5,000 | 04:10:00 | 04:26:39 |
环境与样本生成#
本文固定使用 Python 3.12、uv 0.10.0、PyArrow 25.0.0、parquet-cli 1.16.0 和 DuckDB 1.5.5。先核对环境。PEP 723 会让 uv 按代码块顶部的声明选择 Python 3.12,并安装固定版本的 PyArrow。
uv --version
uv run --python 3.12 python --version
duckdb --version
parquet versionparquet-cli 比查询引擎更适合观察 Schema、Page Header、FileMetaData 和索引。它没有独立的官方二进制发行包,本文按 1.16.0 tag 构建。
构建环境需要 JDK 17 和 Apache Thrift 0.22.0。较新的 JDK 会与该版本依赖的 Hadoop 3.3.0 不兼容。
git clone --depth 1 --branch apache-parquet-1.16.0 \
https://github.com/apache/parquet-java.git
cd parquet-java
./mvnw -pl parquet-cli -am -DskipTests install
./mvnw -pl parquet-cli dependency:copy-dependencies
PARQUET_JAVA_HOME="$(pwd)"
parquet() {
java -cp "$PARQUET_JAVA_HOME/parquet-cli/target/parquet-cli-1.16.0.jar:$PARQUET_JAVA_HOME/parquet-cli/target/dependency/*" \
org.apache.parquet.cli.Main --dollar-zero parquet "$@"
}
parquet version这是官方 README 给出的 no-Hadoop classpath 方式。命令使用普通 jar 和依赖目录,不把 runtime jar 混进 classpath。runtime jar 会重定位 Avro 类,并可能改变方法签名。
下面是本文唯一的样例生成器。正文引用的文件大小、偏移和查询结果都来自它的默认输出。
展开 generate_parquet_samples.py
# /// script
# requires-python = ">=3.12,<3.13"
# dependencies = ["pyarrow==25.0.0"]
# ///
from __future__ import annotations
import argparse
import hashlib
from collections.abc import Callable, Sequence
from contextlib import suppress
from datetime import datetime, timedelta
from pathlib import Path
import pyarrow as pa
import pyarrow.parquet as pq
SEED = 20260806
SMALL_ROWS = 30_000
SMALL_ROW_GROUP_ROWS = 10_000
_BASE_ID = 1_259_328_812_206_792_704
_BASE_USER_ID = 102_117_463_212
_BASE_REFERENCE_TIME = 1_778_205_786
_BASE_EVENT_TIME = datetime(2026, 5, 8, 0, 0, 0)
_ROWS_PER_USER = 20
_NULL_TIMESTAMP_OFFSET = 500
_TYPE_CODES = (1085, 1005, 1, 1032)
_STAKE_VALUES = (1_000, 15_000, 100_000, 250_000)
_MERCHANT_IDS = (102, 205, 319)
def transaction_schema() -> pa.Schema:
return pa.schema(
[
pa.field("id", pa.int64(), nullable=False),
pa.field("mch_id", pa.int32(), nullable=False),
pa.field("user_id", pa.int64(), nullable=False),
pa.field("out_trans_no", pa.string(), nullable=False),
pa.field("amount", pa.int64(), nullable=False),
pa.field("balance", pa.int64(), nullable=False),
pa.field("created_at", pa.timestamp("us"), nullable=True),
pa.field("updated_at", pa.timestamp("us"), nullable=True),
]
)
def nested_schema() -> pa.Schema:
action = pa.struct(
[
pa.field("type", pa.string(), nullable=False),
pa.field("amount", pa.int64(), nullable=True),
]
)
return pa.schema(
[
pa.field("event_id", pa.int64(), nullable=False),
pa.field(
"user",
pa.struct([pa.field("city", pa.string(), nullable=True)]),
nullable=True,
),
pa.field(
"actions",
pa.list_(pa.field("element", action, nullable=False)),
nullable=True,
),
]
)
def nested_events_table() -> pa.Table:
rows = [
{
"event_id": 1,
"user": {"city": "Shanghai"},
"actions": [
{"type": "bet", "amount": -100},
{"type": "win", "amount": 250},
],
},
{"event_id": 2, "user": None, "actions": []},
{"event_id": 3, "user": {"city": None}, "actions": None},
{
"event_id": 4,
"user": {"city": "Beijing"},
"actions": [{"type": "bet", "amount": None}],
},
]
return pa.Table.from_pylist(rows, schema=nested_schema())
def _financial_value(*, user_index: int, within_user: int, seed: int) -> int:
pair_within_user, pair_position = divmod(within_user, 2)
stake = _STAKE_VALUES[(user_index * 3 + pair_within_user + seed) % len(_STAKE_VALUES)]
payout = 2 + (user_index + pair_within_user + seed) % 4
amount = stake * payout if pair_position == 1 else -stake
return amount
def generate_transaction_batch(
*,
start_row: int,
row_count: int,
seed: int,
) -> pa.RecordBatch:
schema = transaction_schema()
columns: dict[str, list[object]] = {field.name: [] for field in schema}
seed_offset = seed % 997
active_user_index = -1
running_balance = 0
for index in range(start_row, start_row + row_count):
user_index, within_user = divmod(index, _ROWS_PER_USER)
pair_within_user, pair_position = divmod(within_user, 2)
global_pair = user_index * (_ROWS_PER_USER // 2) + pair_within_user
is_win = pair_position == 1
type_code = _TYPE_CODES[(user_index + pair_within_user + seed) % len(_TYPE_CODES)]
amount = _financial_value(
user_index=user_index,
within_user=within_user,
seed=seed,
)
user_id = _BASE_USER_ID + user_index
reference_time = _BASE_REFERENCE_TIME + global_pair * 6
game_instance = 281 + user_index % 5_000
reference_segment = f"{type_code}-0-{game_instance}-{reference_time}-0"
action_marker = 1 if is_win else 0
serial = 17_782_057_860_228 + global_pair * 317 + pair_position
is_null_timestamp = index > 0 and index % 1_000 == _NULL_TIMESTAMP_OFFSET
event_time = (
None
if is_null_timestamp
else _BASE_EVENT_TIME + timedelta(seconds=index)
)
if active_user_index != user_index:
running_balance = 10_000_000_000 + (user_index % 5_000) * 1_000_000
for previous_position in range(within_user):
previous_amount = _financial_value(
user_index=user_index,
within_user=previous_position,
seed=seed,
)
running_balance += previous_amount
active_user_index = user_index
running_balance += amount
columns["id"].append(_BASE_ID + index * 33_554_433 + seed_offset)
columns["mch_id"].append(_MERCHANT_IDS[(user_index // 5_000) % len(_MERCHANT_IDS)])
columns["user_id"].append(user_id)
columns["out_trans_no"].append(
f"{user_id}-{action_marker}{serial}-{reference_segment}"
)
columns["amount"].append(amount)
columns["balance"].append(running_balance)
columns["created_at"].append(event_time)
columns["updated_at"].append(event_time)
arrays = [pa.array(columns[field.name], type=field.type) for field in schema]
return pa.RecordBatch.from_arrays(arrays, schema=schema)
def _atomic_write(
final_path: Path,
writer_profile: str,
write_temporary: Callable[[Path], None],
) -> Path:
final_path.parent.mkdir(parents=True, exist_ok=True)
temporary_path = final_path.with_name(f"{final_path.name}.tmp")
temporary_path.unlink(missing_ok=True)
try:
write_temporary(temporary_path)
temporary_path.replace(final_path)
except Exception as error:
with suppress(OSError):
temporary_path.unlink(missing_ok=True)
raise RuntimeError(
f"failed to write {writer_profile} Parquet file: {final_path}"
) from error
return final_path
def write_transactions_small(output_dir: Path) -> Path:
final_path = output_dir / "transactions-small.parquet"
def write(temporary_path: Path) -> None:
with pq.ParquetWriter(
temporary_path,
transaction_schema(),
version="2.6",
compression="snappy",
use_dictionary=True,
write_statistics=True,
data_page_version="2.0",
use_compliant_nested_type=True,
write_page_index=True,
write_page_checksum=True,
max_rows_per_page=1_000,
bloom_filter_options={
"out_trans_no": {"ndv": SMALL_ROWS, "fpp": 0.01},
},
) as writer:
for start_row in range(0, SMALL_ROWS, SMALL_ROW_GROUP_ROWS):
batch = generate_transaction_batch(
start_row=start_row,
row_count=SMALL_ROW_GROUP_ROWS,
seed=SEED,
)
writer.write_batch(batch, row_group_size=SMALL_ROW_GROUP_ROWS)
return _atomic_write(final_path, "small-profile", write)
def write_nested_events(output_dir: Path) -> Path:
final_path = output_dir / "nested-events.parquet"
def write(temporary_path: Path) -> None:
pq.write_table(
nested_events_table(),
temporary_path,
row_group_size=4,
version="2.6",
compression="snappy",
use_dictionary=True,
write_statistics=True,
data_page_version="2.0",
use_compliant_nested_type=True,
write_page_index=True,
write_page_checksum=True,
max_rows_per_page=2,
)
return _atomic_write(final_path, "nested-profile", write)
def write_default_files(output_dir: Path) -> tuple[Path, Path]:
return write_transactions_small(output_dir), write_nested_events(output_dir)
def file_sha256(path: Path) -> str:
digest = hashlib.sha256()
with path.open("rb") as source:
for block in iter(lambda: source.read(1024 * 1024), b""):
digest.update(block)
return digest.hexdigest()
def print_file(path: Path) -> None:
metadata = pq.ParquetFile(path).metadata
assert metadata is not None
print(
f"{path}: rows={metadata.num_rows} row_groups={metadata.num_row_groups} "
f"bytes={path.stat().st_size} sha256={file_sha256(path)}"
)
def build_parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(description="Generate reproducible Parquet samples")
parser.add_argument("--output-dir", type=Path, default=Path("parquet-lab"))
return parser
def main(argv: Sequence[str] | None = None) -> int:
args = build_parser().parse_args(argv)
for output in write_default_files(args.output_dir):
print_file(output)
return 0
if __name__ == "__main__":
raise SystemExit(main())把展开后的代码保存为 generate_parquet_samples.py,默认命令会在 parquet-lab/ 生成两个很小的结构样例:
uv run generate_parquet_samples.py
shasum -a 256 \
parquet-lab/transactions-small.parquet \
parquet-lab/nested-events.parquet同一环境重复运行会原子替换文件,并得到相同字节:
| 文件 | 行数 | Row Groups | 物理字节 | SHA-256 |
|---|---|---|---|---|
transactions-small.parquet |
30,000 | 3 | 1,489,068 | 62a5dee8139f1e80ad4e95decbde08ba5542578999d37cb0e7695a57efaa07f3 |
nested-events.parquet |
4 | 1 | 2,000 | fefbf3bfe261e10619c7bb1d48ff8ac1b37f01ce227c1461920d261275c4fd49 |
复现清单#
按“环境与样本生成”小节准备环境并生成样例后,可以用下面的命令复查文章中的结构和查询结果:
# 1. 文件尾、FileMetaData、Schema
xxd -g 1 -s -8 parquet-lab/transactions-small.parquet
od -An -j 1489060 -N 4 -tu1 parquet-lab/transactions-small.parquet
parquet footer --raw parquet-lab/transactions-small.parquet
parquet schema --parquet parquet-lab/transactions-small.parquet
# 2. Statistics、Row Group
parquet meta parquet-lab/transactions-small.parquet
parquet bloom-filter -c out_trans_no -v missing-transaction-000 \
parquet-lab/transactions-small.parquet
# 3. Column Chunk、Page Index、Page Header
parquet column-size parquet-lab/transactions-small.parquet
parquet column-index -c created_at parquet-lab/transactions-small.parquet
parquet pages --raw -c user_id,amount,created_at \
parquet-lab/transactions-small.parquet
# 4. 查询结果
duckdb -c "
SELECT user_id, SUM(amount)
FROM read_parquet('parquet-lab/transactions-small.parquet')
WHERE created_at >= TIMESTAMP '2026-05-08 03:00:00'
AND created_at < TIMESTAMP '2026-05-08 04:00:00'
GROUP BY user_id;
"附录 C:Random Access Parquet(RAP)#
Parquet 通常由 Trino、BigQuery 一类分析引擎读取。这些引擎擅长并行扫描和聚合,却不适合按一个 key 在线读取少量记录。用户历史、长尾实体详情或 AI Agent 的上下文检索都要求较低延迟。为一次点查启动分布式任务,还会增加调度和查询规划开销。
分区裁剪可以缩小候选文件范围,Bloom Filter 也能排除部分 Row Group。但它们不能直接回答某个 key 位于文件中的哪一行。
读取器进入候选文件后,仍要依次读取 FileMetaData、解析 Row Group 元数据、扫描 key 列,再定位 value Page。每一步都依赖上一步。文件位于对象存储时,多轮 Range Read 的往返延迟会逐次累加。
Random Access Parquet(RAP)针对的正是这段文件内发现过程。它不是新的 Parquet 规范,基础方案也不复制整份 value 数据,而是在现有 Parquet 文件之外增加一条精确访问路径。
用外部索引替代文件内查找#
RAP 的外部索引是一个 multimap,也就是“一键对应多个位置”的映射。同一个 key 可能出现在多个时间分区和文件中,因此不能只保存一个地址。
每条索引记录包含查询 key、Parquet 文件标识、该 key 的行号集合,以及可选的 value count。文件标识可以压缩成字典序号。value count 则让服务端在读取数据前处理分页。
索引由离线任务构建。构建器读取现有文件的 FileMetaData 和目标列的 Page 位置,再扫描 key 列,生成 key -> file + row numbers 映射。
新一批 Parquet 文件到达时,只需追加对应的索引分片,不必重写旧文件。因此,RAP 可以直接用于没有特殊准备的标准 Parquet。
Index Builder 的输出:追加式 multimap#
把 key -> file + row numbers 展开后,Index Builder 的核心产物可以用下面的逻辑模型表示:
FileDictionary: file_ordinal -> Parquet URI
Bucket: hash(key) -> [Fragment]
Fragment: key -> [{file_ordinal, row_numbers[], value_count?}]文件路径通常很长,索引条目因此只保存字典编码后的 file_ordinal。索引较大时,可以按 key 做 hash 分桶;每轮离线任务只为新到达的 Parquet 追加 fragment,旧分片保持不变。
以图中的 user-123 为例,第一个 fragment 指向 file 17 的第 120、121 行。下一批数据又追加了 file 42 的第 8 行。查询时,读取器合并目标 bucket 中的条目,得到 {17: [120, 121], 42: [8]}。
结果仍然只是文件序号和行号。各列对应的 Page 或字节范围,要在下一步通过缓存的 Parquet 文件元数据换算。
Spotify 原文明确提到 multimap 条目、字典编码的文件序号、追加 fragment,以及按 hash 对大索引分桶。图中的文件字典、bucket b7 和日期分片只是说明这些关系的参考模型,不是 RAP 的固定磁盘格式。原文也没有公开 fragment 编码、hash 函数或具体存储引擎。
在线读取分为四步:
- 用 key 查询外部索引,得到一组文件与行号。
- 用缓存的文件元数据把行号换成各目标列的 Page 或字节范围。
- 合并相邻范围,并行发出互不依赖的 Range Read。
- 解压、解码并抽取目标行,最后合并多个文件中的结果。
这里的 O(1) 只指哈希索引按 key 查找,不代表端到端查询复杂度。实际开销仍取决于命中文件数、读取列数、Page 大小、网络请求和解码工作。
Parquet 的 Statistics、Bloom Filter 和 Page Index 用于排除候选范围。RAP 索引返回精确文件与行号,两者职责不同。
适用边界与准备型优化#
外部索引解决了“数据在哪里”,却不一定解决读取放大。标准 Parquet 仍以 Page 作为独立解压边界:即使已经知道目标行,读取器也可能为了很小的结果读取整个 Page。
更看重点查延迟或读取成本时,可以在写入阶段集中相同 key 的数据,例如按 key 排序,或者在 key 边界刷新 Page。也可以把在线服务需要的字段收进 Blob、Protobuf 或 Variant 列,减少每次请求涉及的列数。很小的值和预聚合结果还可以直接放进覆盖索引,省去对象存储读取。
这些做法会增加 Page Index 或外部索引的体积,也可能削弱编码效率和列裁剪效果。是否采用,要结合分析查询的需求衡量。
RAP 适合历史数据、长尾实体以及读多写少的在线点查,也可以作为 AI Agent 获取上下文的入口。它不提供事务、高频行级更新、约束或通用查询执行能力。批量分析仍沿用 Parquet 原有的扫描路径。
参考资料#
- Apache Parquet Overview
- Apache Incubator:Parquet Project Incubation Status
- Google Research:Dremel: Interactive Analysis of Web-Scale Datasets
- Parquet Format 2.13.0 README:File Format 与 Nested Encoding
- Parquet Format 2.13.0
parquet.thrift - Parquet Format 2.13.0 Logical Types
- Parquet Format 2.13.0 Encodings
- Parquet Format 2.13.0 Compression
- Parquet Format 2.13.0 Encryption
- Parquet Format 2.13.0 Page Index
- Parquet Format 2.13.0 Bloom Filter
- PyArrow 25.0.0
write_table - parquet-cli 1.16.0 README
- parquet-cli 1.16.0
Util.java:Codec 与 Encoding 缩写 - parquet-cli 1.16.0
ShowPagesCommand.java:Page 输出列 - DuckDB:Reading and Writing Parquet Files
- DuckDB 1.5.5:Parquet Row Group 裁剪源码
- Spotify Engineering:Indexing the Data Lake for Online Point Queries