Apache Parquet 文件格式深度实战:从 Dremel 记录拆解到 Page Index 谓词下推的工程全链路
如果你做过数据湖,一定听过一句话:"Parquet 是列存的。"这句话对,但它掩盖了 Parquet 真正难的地方。
列存本身不难——把同列的值堆在一起写就行了。Parquet 真正的工程价值在于:它能在保持列式压缩效率的同时,无损地表达任意深度的嵌套结构(struct / list / map),并且让查询引擎只读它需要的那几个字节。前者靠 Google Dremel 论文的记录拆解算法,后者靠 Footer 元数据 + Page Index 的多级裁剪。这两件事,才是 Parquet 能在十年后依然是数据湖事实标准的原因。
本文从字节布局开始,一路拆到记录装配算法和谓词下推,并给出可直接运行的代码。
一、文件物理布局:四个 "PAR1" 字节决定的世界
一个 .parquet 文件的骨架是这样的:
┌──────────────────┐
│ PAR1 │ 4 字节 magic(文件头)
├──────────────────┤
│ Row Group 0 │
│ ├ Column Chunk a│ (Data Page V2 × N + Dictionary Page)
│ ├ Column Chunk b│
│ └ Column Chunk c│
├──────────────────┤
│ Row Group 1 │
│ └ ... │
├──────────────────┤
│ Column Index │ 页级 min/max(可选,Parquet 1.11+)
│ Offset Index │ 页级 offset / row boundary(可选)
│ Bloom Filter │ Split Block Bloom Filter(可选)
├──────────────────┤
│ Footer │ Thrift Compact 编码的 FileMetaData
│ <footer_len> │ 4 字节小端
│ PAR1 │ 4 字节 magic(文件尾)
└──────────────────┘
注意收尾那 8 个字节。读 Parquet 不是从头读,而是从尾读:seek 到 filesize - 8,读出 footer 长度,再 seek 到 footer 起点反序列化 Thrift 结构。这就是为什么 Parquet 对对象存储(S3 / OSS)友好——一次小范围 range GET 就能拿到整个文件的"目录",不需要把几百 MB 全拉下来。
层级关系要记牢:
- Row Group:一组行的横向切片,是并行度和裁剪粒度的基本单位。parquet-mr 默认目标 128MB。
- Column Chunk:Row Group 内某一列的全部数据,物理上连续。
- Page:Column Chunk 内的压缩/编码单元,默认 1MB(data page)。压缩和编码都发生在 Page 这一层,不是整个 chunk。
这个"Row Group → Column Chunk → Page"的三级,正好对应后面裁剪的三个层次:Row Group 裁剪 → 列裁剪 → Page 裁剪。
二、嵌套结构的核心难题:Dremel 记录拆解
列存面对嵌套数据的困难是:JSON 那样的树,a.b.c 可能一行里有 0 个、1 个或 100 个值。简单拍平会丢失"哪个值属于哪一层"的信息。
Parquet 的解法是给每个值附加两个整数:Repetition Level (r) 和 Definition Level (d)。
- d(定义层):从根到该字段的路径上,有多少个 *optional/repeated* 字段实际存在。
d == max_d表示值非空;d < max_d表示值为 NULL(并指明 NULL 发生在哪一层)。 - r(重复层):该值是在路径的哪一级"重复"产生的。
r == 0意味着新纪录开始。
看一个具体的电商事件 schema:
message Event {
required int64 ts;
required string user_id;
repeated group Items { // r 层 = 1, d 层 = 1
required string sku; // required → max_d = 1
optional double price; // optional → max_d = 2
repeated string tags; // r 层 = 2, max_d = 2
}
}
两条记录:
- 记录 A:
Items = [{sku:'A1', price:9.9, tags:['sale','new']}, {sku:'A2', price:None, tags:[]}] - 记录 B:
Items = [](无 items)
各列的 (r, d) 流:
| 列 | 值 | r | d | 含义 |
|---|---|---|---|---|
Items.sku | 'A1' | 0 | 1 | 新纪录,Items 存在 → required 必有值 |
| 'A2' | 1 | 1 | Items 重复了一次 | |
| null | 0 | 0 | 新纪录,但 Items 本身缺失 | |
Items.price | 9.9 | 0 | 2 | 新纪录,price 有值 |
| null | 1 | 2 → 1 | Items 重复,但 price 缺失(NULL 发生在 optional 层) | |
| null | 0 | 0 | 新纪录,Items 缺失 | |
Items.tags | 'sale' | 0 | 2 | 新纪录,tags 有值 |
| 'new' | 2 | 2 | tags 这一级重复(非 Items 重复) | |
| null | 1 | 1 | Items 重复,但 tags 是空列表 | |
| null | 0 | 0 | 新纪录,Items 缺失 |
关键点在 Items.tags 的 'new':r=2 说明重复发生在 tags 层级;而 sku 的 'A2' 是 r=1,说明重复发生在 Items 层级。正是这个区分,让装配器能还原出精确的树形。
注意 Items.price 那行我特意标注了 2 → 1:缺失的 optional 字段不消耗值槽位,但消耗一个 (r,d) 条目。这是装配算法能对齐的前提。
三、把 (r, d) 流装回树:可运行的装配代码
下面是完整可跑的装配实现,输入是各列的 (values, defs, reps):
def split_records(defs, reps):
"""按 r=0 把 (d, r) 流切分成 per-record 的 definition 列表"""
out = []
for d, r in zip(defs, reps):
if r == 0 or not out:
out.append([])
out[-1].append(d)
return out
def assemble(sku, price, tags):
"""sku=(values, defs, reps),其余同。返回 record 列表"""
sku_recs = split_records(sku[1], sku[2])
price_recs = split_records(price[1], price[2])
tag_recs = split_records(tags[1], tags[2])
records = []
sv, pv, tv = iter(sku[0]), iter(price[0]), iter(tags[0])
for s_rec, p_rec, t_rec in zip(sku_recs, price_recs, tag_recs):
# 1) 用 required 锚列建 item 骨架(d==1 表示 Items 存在)
items = []
for d_s in s_rec:
if d_s < 1:
continue
items.append({'sku': next(sv), 'price': None, 'tags': []})
# 2) 回填 optional:只有 d == max_d(2) 才消耗一个值
p_iter = iter(p_rec)
for it in items:
if next(p_iter, 0) == 2:
it['price'] = next(pv)
# 3) 回填 repeated:d==2 连续消耗,直到 d<2 表示该 item 的 list 结束
t_iter = iter(t_rec)
for it in items:
d = next(t_iter, 0)
while d == 2:
it['tags'].append(next(tv))
d = next(t_iter, 0)
records.append({'items': items})
return records
# ── 用第二节的两条记录验证 ──────────────────────────────
sku = (['A1', 'A2'], [1, 1, 0], [0, 1, 0])
price = ([9.9], [2, 1, 0], [0, 1, 0])
tags = (['sale', 'new'], [2, 2, 1, 0], [0, 2, 1, 0])
import json
print(json.dumps(assemble(sku, price, tags), ensure_ascii=False))
输出:
[{"items": [{"sku": "A1", "price": 9.9, "tags": ["sale", "new"]},
{"sku": "A2", "price": null, "tags": []}]},
{"items": []}]
与原始输入完全一致。这段代码揭示了一个重要工程事实:装配一棵嵌套树不需要随机访问,只需要三个线性游标。这就是 Parquet 能做到流式、低内存的原因,也是 DuckDB / Velox 这类引擎能向量化地批量装配的基础。
四、编码体系:为什么 Parquet 比"列存 + gzip"强得多
Parquet 的压缩是编码 + 通用压缩两级流水线。很多人只配了 compression='snappy' 就以为完事了,实际上真正省空间的是编码层。
| 编码 | 适用 | 原理与收益 |
|---|---|---|
PLAIN | 兜底 | 裸值,几乎不省 |
RLE_DICTIONARY | 低基数字符串/枚举 | 字典把长字符串压成 int32 索引,再对索引 RLE。基数是决定性因素 |
RLE + BIT_PACKED 混合 | definition/repetition levels | 见下 |
DELTA_BINARY_PACKED | 单调递增 int(时间戳、自增 ID) | 存差值 + 每个值的位宽,比字典更省 |
DELTA_BYTE_ARRAY | 排序后的字符串(URL、前缀相似) | 存公共前缀长度 + 后缀 |
BYTE_STREAM_SPLIT | float / double | 把 8 字节拆成 8 个定长流,让下游通用压缩率暴涨 |
BYTE_STREAM_SPLIT 值得单独说:浮点数的尾数近乎随机,gzip/zstd 压不动。拆成 8 条字节流后,每条流内部高度规整(比如指数位大量相同),配合 ZSTD 实测能把 float64 再压掉 20%–40%。这是 Parquet 独有的、很多工程师不知道的优化点。
definition/repetition levels 用的是 RLE / Bit-Packing Hybrid,格式如下:一个 varint 头 (len << 1) | flag,flag=0 是 RLE 重复 run,flag=1 是 bit-packed 组(每组 8 个值,LSB 优先)。解码实现:
def read_varint(buf, pos):
result = shift = 0
while True:
b = buf[pos]; pos += 1
result |= (b & 0x7F) << shift
if not (b & 0x80):
return result, pos
shift += 7
def rle_bitpack_decode(buf, bit_width, count):
"""Parquet RLE / Bit-Packing Hybrid 解码,返回 count 个非负整数"""
out, pos = [], 0
while len(out) < count:
header, pos = read_varint(buf, pos)
if header & 1: # bit-packed
ngroups = header >> 1
nbytes = ngroups * bit_width
bits = int.from_bytes(buf[pos:pos + nbytes], 'little')
pos += nbytes
mask = (1 << bit_width) - 1
for i in range(ngroups * 8):
out.append((bits >> (i * bit_width)) & mask)
else: # RLE run
run = header >> 1
nbytes = (bit_width + 7) // 8
val = int.from_bytes(buf[pos:pos + nbytes], 'little')
pos += nbytes
out.extend([val] * run)
return out[:count]
bit_width 由 schema 最大层数决定(通常 ≤ 3,即 2 bit)。一个百万行的文件,levels 数据常常只有几十 KB。
五、谓词下推的三级火箭
这才是 Parquet 在查询侧真正的杀手锏。一次 WHERE ts BETWEEN x AND y AND sku = 'A1',引擎最多可以做三层裁剪:
第一级:Row Group 裁剪。 Footer 里每个 Column Chunk 都带 Statistics(min/max、null_count、distinct_count)。引擎先读 Footer,用 min/max 判断哪些 Row Group 完全不可能命中,直接跳过。
第二级:Page 裁剪(Column Index)。 Parquet 1.11 引入的 Page Index 把统计信息下沉到 Page 粒度:ColumnIndex 存每个 page 的 min/max/null_count,OffsetIndex 存每个 page 在文件中的 offset 和起始行号。这样即使某个 Row Group 命中,也只需要读其中少数几个 page。
第三级:Bloom Filter。 对高基数列的等值/IN 查询,min/max 几乎无效(取值范围太宽)。Parquet 的 Split Block Bloom Filter 能给出"该 Row Group 绝对不含此值"的确定性否定,把随机 IO 降到接近零。
用 pyarrow 观察这些元数据:
import pyarrow.parquet as pq
pf = pq.ParquetFile("events.parquet")
md = pf.metadata
print("row groups:", md.num_row_groups, "rows:", md.num_rows)
rg = md.row_group(0)
for i in range(rg.num_columns):
col = rg.column(i)
st = col.statistics
print(f"{col.path_in_schema:20s} "
f"min={st.min} max={st.max} nulls={st.null_count} "
f"compressed={col.total_compressed_size}")
# Page Index 需要显式读取
ci = pf.metadata.row_group(0).column(0).statistics
print("has stats:", ci is not None, "| distinct:", ci.distinct_count)
生产经验:Row Group 大小是个真 trade-off。太大 → 裁剪粒度粗,且单个 Row Group 必须整体进内存(列式读是 chunk 粒度);太小 → Footer 元数据膨胀,S3 请求数增加,且每个 chunk 的压缩字典无法共享。128MB~512MB 是常见甜区;时序/高过滤率场景可以压到 64MB 换取更细的裁剪。
六、生产环境踩过的坑
1. 小文件问题(Small File Problem)是 Parquet 的头号杀手。 每个文件都有独立的 Footer 和独立的字典,1 万个 1MB 的文件在查询侧比 10 个 1GB 的文件慢一个数量级——不是 IO 慢,是元数据解析和 S3 GET 请求数爆炸。必须用 compaction(后台合并)兜底。
2. 不要在写入时忽略排序。 sort_columns 对压缩率的影响远超压缩算法选择。按 user_id, ts 排序后,字典编码基数骤降、delta 编码差值变小,实测体积能再降 30%–50%。代价是写入端需要 buffer 整个 Row Group。
3. Schema 演进要显式打 field_id。 如果上层用了 Iceberg / Delta 这类表格式,务必通过 parquet.field.id 给每个字段写死 ID。否则重命名列后,按位置匹配的引擎会静默读错数据——这类 bug 极难排查。
4. 压缩算法别默认选 gzip。 列式数据解压常在数据扫描热路径上:ZSTD 在相近压缩率下解压速度是 gzip 的 3–5 倍,SNAPPY 解压最快但压缩率明显差。通用建议是 ZSTD level 3,除非存储成本完全不敏感。
5. 嵌套别超过必要深度。 每多一层 repeated,definition/repetition levels 的 bit_width 和处理开销就增加一分,装配时的分支预测失败也会上升。能用 struct 表达的就别套 list<struct<list<...>>>。
七、结语
Parquet 常被简化成"列式存储格式",但真正让它在数据湖里不可替代的是三件事:Dremel 记录拆解让嵌套结构无损列存化、编码体系针对列内数据分布做了深度定制、Footer + Page Index + Bloom Filter 构成了从文件到字节的三级裁剪。
理解第二、三节,你就能自己写一个 Parquet reader;理解第五、六节,你才能把 Parquet 在生产环境里用对。大多数 Parquet 性能问题,根因不是"Parquet 慢",而是 row group 太大、文件太小、或者忘了排序。

发表评论 取消回复