Parquet 文件格式解析

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 可得到结果列、过滤条件和聚合方式,却无法告诉读取器这些列数据在文件中的具体位置。执行查询前还要回答三个问题:

  1. Schema 如何定义,user_idamountcreated_at 对应哪些叶子列?
  2. 哪些 Row Group 和 Page 不可能满足时间条件,可以直接跳过?
  3. 候选列采用什么编码和压缩,读取器怎样找到并解码它们?

1. Parquet 文件结构#

先看完整文件,再模拟查询。图中的粗边界表示 File -> Row Group -> Column Chunk -> Page。细线只表示局部放大,不是数据流。

这几层各有分工。Row Group 划分一批行,Column Chunk 保存这批行中的一个叶子列,Page 再保存该列的一小段编码数据。每层的大小都由写入器选择,不是格式常量。

Bloom Filter 和 Page Index 不在这条包含链中。Bloom Filter 用来快速判断“某个值一定不存在”。Page Index 由 ColumnIndexOffsetIndex 组成:前者像页级范围目录,记录每页的边界;后者像位置表,记录页所在的字节范围和行区间。

Parquet 文件包含多个 Row Group;图中展开 Row Group 1 的 Column 1、Column 2 到 Column 8,再展开 Column 1 的可选 Dictionary Page 和多个 Data Page;文件尾的 FileMetaData、4 字节长度和末尾 PAR1 用于定位这些结构
Parquet 的文件层级

这个样例写入器的物理布局如下;它描述的是本文件的写入顺序,而不是规范强制的可选辅助结构排列顺序。

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、ColumnIndexOffsetIndex 都服务于单个 Column Chunk,但字节不在该 Chunk 的连续 Page 区域内。读取器要先解析文件尾的 FileMetaData,再根据其中的引用找到它们。因此,这张图描述文件层级,不描述读取顺序。

2. 用 SQL 驱动一次读取#

查询引擎解析开篇 SQL 后,会把投影和过滤条件交给 Parquet 扫描算子。投影需要 user_idamount,过滤需要 created_at,因此扫描阶段要读三个叶子列。

GROUP BY user_idSUM(amount) 由上层执行算子完成。Parquet 读取器只负责返回通过时间条件的列数据,不负责聚合。

未加密 Parquet 文件的读取路径:读取器从文件尾解析 FileMetaData,根据 Schema 选择投影列和谓词列,以 Statistics、Bloom Filter 以及实现支持时的 Page Index 排除候选范围,解码后执行剩余谓词
Parquet 读取路径

先区分格式能力与工具的实际行为。样例生成器会写出 Page Index,parquet-cli 也能解析它。不过,PyArrow 25.0.0 文档明确说明,其读取端尚不使用 Page Index。

DuckDB 1.5.5 的 Parquet 扫描器会用 Row Group Statistics 和 Bloom Filter 做裁剪。它没有读取 column_index_offsetoffset_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_idSUM(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#

文件最后依次保存 4186 字节 FileMetaData、4 字节小端 FileMetaData length 和 4 字节 PAR1;读取器按 file_metadata_start = file_size - 8 - file_metadata_length 计算 FileMetaData 起点
从文件尾定位 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_length

FileMetaData 不直接装入所有辅助结构。它主要保存 Schema、Row Group、Column Chunk,以及通往其他字节区域的引用。读取器解析这些引用后,才会跳到具体的 Page 或索引。

引用分属两个层次:

所属结构 保存的内容
ColumnChunk ColumnIndexOffsetIndex 的 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.parquet
message 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_idamount,使用 created_at 谓词过滤,因此扫描算子只需继续追踪这三条叶子路径即可。

Schema 还规定字段是否可空,以及它的物理表示和逻辑语义。created_at 的 Physical Type 是 INT64,Logical Type 是 TIMESTAMP(MICROS,false)。读取器因此把整数解释为微秒精度、无 UTC 调整语义的时间戳。out_trans_no 的 Physical Type 是 BYTE_ARRAY,其 STRING 注解表示这些字节采用 UTF-8。

Schema 将 created_at、out_trans_no 和 mch_id 的字段语义映射到 INT64、BYTE_ARRAY 与 INT32 等叶子物理类型
Schema 到叶子物理类型的映射

扫描算子现在已经确定 created_atuser_idamount 三条叶子路径。下一步先看 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.parquet
Row 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:00

RG0 的最大值早于查询下界,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.parquet
Row 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 可以独立读取。

30,000 条完整记录按写入顺序横向切成三个 Row Group,每组 10,000 行,并各自保留相同的八个叶子列
完整记录按行切成三个 Row Group

parquet meta 报告 RG1 有 10,000 条记录,压缩后的全部 Column Chunk 合计约 464.092 kB:

parquet meta parquet-lab/transactions-small.parquet
Row 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_idamountcreated_at。其余五列可以跳过。

逻辑上每个 Row Group 都有 8 个叶子 Column Chunk;右侧只摘录 RG0 物理顺序中的第 1、3、4、7 列,说明同一 Row Group 内各 Column Chunk 依次连续存放
Column Chunk 的逻辑与物理布局

同一条 parquet meta 命令会列出 RG1 的八个 Column Chunk。下面只保留本次扫描涉及的三行:

parquet meta parquet-lab/transactions-small.parquet
Row 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 的压缩后大小差异很大,但读取器不必靠顺序扫描猜测边界。它会分别使用 dictPageOffsetdataPageOffsettotalSize 定位。本文样例的具体值集中列在附录 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

一个 Column Chunk 由可选的 Dictionary Page 和连续 Data Page 组成;每张 Page 都使用变长 PageHeader 加 Page body 的通用结构,type 决定 Dictionary、Data V1 或 Data V2 专属 Header 和页体布局
Column Chunk 中的 Page 序列、通用 PageHeader 与 type 分支

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 指 ColumnIndexOffsetIndex 两个独立结构,它们的位置和长度由 ColumnChunk 元数据引用。

CRC 字段位于 PageHeader 内,校验范围却只覆盖写入磁盘的 Page body,不包含 Header。读取器先反序列化通用 Header,再根据 type 解释页专属 Header 和后续 compressed_page_size 字节。

7.2 Page Index:ColumnIndexOffsetIndex#

格式层已经把范围缩小到 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 并没有执行这一步。

支持 Page Index 的读取器可以先用 Column Chunk Statistics 排除 RG0 和 RG2,再用可选 Page Index 将 RG1 缩小到 Data Page 0 至 Page 4;Bloom Filter 只为等值条件提供确定不存在或可能存在的判断
支持 Page Index 的读取器可执行两级裁剪

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_idamount 中覆盖相同行区间的 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 编码也不代表页体一定经过压缩。

Parquet 先按数据形状选择 PLAIN、RLE 或 Dictionary 编码,再把编码后的 Page 页体交给压缩 Codec
编码与 Codec 压缩
编码 保存什么 适用场景
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.parquet
page   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 字段树;中间的虚线表示引用,不是数据流。

Parquet 完整文件结构:左侧展开 Row Group 0、created_at Column Chunk 和 Data Page 0,右侧展开 FileMetaData、RowGroup、ColumnChunk 与 ColumnMetaData,中间引用线连接 Bloom Filter、ColumnIndex 和 OffsetIndex
物理数据与 FileMetaData 引用

样例文件的精确偏移图和对应数值放在附录 B,避免在主线中反复穿插字节验算。

开篇 SQL 对这些结构的使用顺序可以压缩成一条完整链路:

  1. 读取末尾 8 字节,根据长度字段定位 FileMetaData。
  2. 从 Schema 解析 created_atuser_idamount,其余五个叶子列不进入扫描。
  3. created_at Statistics 排除 RG0、RG2,保留 RG1。
  4. 沿 RG1 的三个 Column Chunk 引用跳到各自的 Page 区域。
  5. 若读取器实现 Page Index 裁剪,可用 ColumnIndex 保留 created_at 的 Data Page 0 至 Page 4。随后用 OffsetIndex 取得字节范围和行区间。本文的 DuckDB 1.5.5 查询没有走这条路径。
  6. 读取器校验、解压、解码候选页,执行剩余时间过滤,再把 user_idamount 交给聚合算子。

物理元数据只能说明哪些数据可以读取或跳过。最后还要用独立实现核对查询结果。下面的 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.parquet
message 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.amount

DL 表示从根走到叶子时,可选或重复节点定义到了多深。RL 表示当前条目与前一个条目共享到哪一层重复路径。

null、empty list 等情况仍会产生 level 条目,但不会向物理 values 区写入占位值。按上面四条记录展开,得到:

四条嵌套记录被拆成三个叶子流;每个流分别列出 level 条目数、真正写入的非 null 物理值以及 DL 和 RL,说明 null、空列表和 null 列表只占 level 条目,不会写入空值占位符
嵌套数据的 DL/RL 编码
叶子列 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.parquet
Page 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:可复现实验与检查清单#

样例字节与偏移明细#

正文只保留结构与查询主线。本节集中列出样例文件的偏移、长度和校验值,供复现或排错时查阅。

样例文件的局部偏移图,列出 RG0 Column Chunk、RG1 Dictionary 与 Data Page,以及 Bloom Filter、ColumnIndex 和 OffsetIndex 的真实位置
样例文件的关键字节偏移

文件尾的定位数据如下:

项目 计算或含义
文件大小 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_sizeuncompressed_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 version

parquet-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 文件之外增加一条精确访问路径。

标准 Parquet 点查需要依次读取 FileMetaData、解析 Row Group 元数据、扫描 key 列并定位 value Page;RAP 通过外部精确索引和缓存的 Page 映射,直接并行读取目标字节范围
标准 Parquet 的依赖读取链与 RAP 的索引读取路径

用外部索引替代文件内查找#

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,旧分片保持不变。

RAP 索引的参考逻辑模型:文件字典把文件序号映射到 Parquet 路径,同一 hash bucket 下的追加分片分别保存 key 对应的文件序号、行号集合和可选计数,查询时合并这些位置后再通过元数据缓存定位 Page
根据 Spotify 原文抽象的索引逻辑模型

以图中的 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 函数或具体存储引擎。

在线读取分为四步:

  1. 用 key 查询外部索引,得到一组文件与行号。
  2. 用缓存的文件元数据把行号换成各目标列的 Page 或字节范围。
  3. 合并相邻范围,并行发出互不依赖的 Range Read。
  4. 解压、解码并抽取目标行,最后合并多个文件中的结果。

这里的 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 原有的扫描路径。

参考资料#