首页
/ Polars Parquet 读取器源码解析:嵌套类型反序列化与 Definition Levels 位图等价优化

Polars Parquet 读取器源码解析:嵌套类型反序列化与 Definition Levels 位图等价优化

2026-09-05 11:21:26作者:农烁颖Land

Polars 的 Parquet 列式读取器位于 polars-parquet 晶格中的 arrow/read 模块,其设计笔记(crates/polars-parquet/src/arrow/read/README.md)揭示了两个核心工程结论:当最大 repetition level 为 0 且最大 definition level 为 1 时,RLE 编码的 definition levels 与 Arrow 的 validity 位图在 LSB 位序上完全等价,可以直接 memcpy;而嵌套 Parquet 组的行结构则编码在 repetition/definition levels 中,需要以递归方式重建 Arrow 的 list offsets 与 struct validity。本文基于该设计文档,结合 nested_utils.rsnested.rs 等源码,完整讲清嵌套 Parquet 数据是如何逐层解码为 Arrow 数组的。

核心观察一:Definition Levels 与 Arrow 位图的 LSB 等价

设计文档(README.md)的第一条 Observations 指出:

当最大 repetition level 为 0 且最大 definition level 为 1 时,RLE 编码的 definition levels 与 Arrow 的 bitmap 精确对应,可以直接 memcpy,无需任何额外转换。

这正是「扁平可选列」(顶层 optional 列,如一个可空 i64 列)的典型场景:每个值只有两种 definition level——0(null)或 1(非 null)。此时 def levels 就是一个 bit-per-value 的位串,且 Parquet Hybrid RLE 的 bitpacking 采用 LSB first 位序(参见 hybrid_rle/bitmap.rs 中 "Writes an iterator of bools into writer, with LSB first" 的写入约定),与 Arrow bitmap 的位序一致,因此整块拷贝即可得到 validity 位图。

源码印证了这一快路径:在 utils/mod.rs 中:

  • page_validity_decoder(L114-L118)直接用 HybridRleDecoder::new(validity, 1, page.num_values())1 bit width 解码 def levels 缓冲区作为页级 validity;
  • State::new(L44-L64)在列是 Optional 且页内确实有 null 时才构建 page_validitynull_count == 0 的页则直接跳过位图解码;
  • 非嵌套解码主循环 unspecialized_decode(L120-L282)中,validity 通过 validity.extend_from_bitmap(&page_validity) 一步扩展(L207),null 位置仅填入该物理类型的默认值 T::default(),不触发真实解码——这就是"def levels 即位图"优化的落地形态。

核心观察二:嵌套 Parquet 组的递归反序列化

设计文档第二节描述了嵌套读取的完整机制,这也是本文的重点。

levels 编码了嵌套结构

Parquet 用 Dremel 模型将嵌套行编码为 primitive 叶子列的 repetition levels(rep)与 definition levels(def)两列。文档指出:"嵌套 parquet 组的行编码在 repetition 和 definition levels 中。在 Arrow 中,它们对应:list 的 offsets 和 validity、struct 的 validity。"

也就是说,Arrow 侧用来表达嵌套形态的三种产物——ListArray 的 offsets 数组、ListArray 的 validity 位图、StructArray 的 validity 位图——在 Parquet 侧并没有独立存储,而是全部可以从 (rep, def) 序列推导出来。

两阶段算法:先收集 nested_info,再递归弹出

文档给出的算法分两个阶段:

  1. 收集阶段:嵌套字段先被递归遍历,收集每一层"是 Struct 还是 List、required 还是 optional",存入 nested_info: Vec<Box<dyn Nested>>(在当前仓库中对应具体类型 Nested,非 trait object)。Nested 根据类型与可空性接收 def/rep levels,并据此累积每层的 offsets 与 validity。
  2. 弹出阶段:一个字段读完之后,从 nested_info 中递归地 pop,逐层构建 StructArrayListArray

文档最后总结出嵌套路径相对扁平路径仅有的两点差异,这两点在源码中都有明确对应:

  1. 不再利用 bitmap 优化,需要把 repetition 与 definition levels 反序列化为整数值;
  2. definition levels 被反序列化两次:一次用于扩展值/可空性,一次用于扩展 nested_info

下面结合源码逐条展开。

Nested / NestedState:逐层累积 offsets 与 validity

nested_utils.rs 中的 Nested 结构体(L12-L22)持有:

  • validity: Option<BitmapBuilder>:仅当该层可空时才存在,用 Option 的有无直接表达可空性;
  • content: NestedContent:枚举 Primitive / List { offsets } / FixedSizeList { width } / Struct(L24-L30),List 分支内嵌 Vec<i64> 偏移数组——这就是 Arrow list offsets 的直接来源;
  • num_valids / num_invalids 两个批量计数器(L17-L21),注释明确说明其目的:"We batch the collection of valids and invalids to amortize the costs",即把连续的 valid/invalid 段累积成 count,遇到交替时才真正调用 extend_constant,从而摊薄位图写入开销。

NestedState::pop(L333-L335)返回 (len, offsets, validity) 三元组,正是文档所说"finish a field 后递归弹出"的实现:nested.rs 中 struct 分支在 L96 执行 nested.pop() 取出 struct 层的 length 与 validity,随后用 StructArray::new(L125-L131)组装各字段数组。

NestedState::levels(L348-L368)则从每层的 is_nullable() / is_repeated() 属性累加出每一深度的 def/rep 阈值,例如 def_delta = is_nullable + is_repeatedrep_delta = is_repeated,并预填初始值 0。这组阈值是后续逐值判断"本值在 depth 是否已定义"(def >= def_levels[depth])和"是否属于新行"(rep == 0)的依据。

collect_nested:把 rep/def 解码为整数并做两遍消费

文档说的"反序列化为 i32"与"def levels 消费两次"体现在 PageDecoder::collect_nestednested_utils.rs):

  1. 每页调用 level_iters(L564-L576):从页数据中 split_buffer 拆出 def/rep 缓冲区,用列描述符中的 max_def_level / max_rep_level 计算 bit width,构建两个 HybridRleDecoder
  2. collect_level_values(L371-L391)把 Hybrid RLE 解码结果物化为 Vec<u16>——这就是文档第 1 点:不再走"LSB 位图直接 memcpy"快路径,而是把 levels 完整展开为逐值整数;
  3. decode_nested(L409-L561)逐值消费:对每一对 (def, rep),沿深度循环判断该层是否已定义、是否有效,调用 nest.push(length, is_valid) 累积各层的 offsets 与 validity 计数——这是 def levels 的第一次消费(扩展 nested_info);同时通过 BatchedNestedDecoder(L214-L284)生成叶子值的 leaf_filterleaf_validity 两个位图——这是第二次消费(扩展值/可空性);
  4. 最后 utils::State::new_nested + state.decode(L672-L683)用叶子 validity/filter 真正解码 primitive 值;
  5. 解码完成后执行 nested_state.pop() 弹出 Primitive 层(L689-L693),留下待上层 columns_to_iter_recursive 继续弹出 List/Struct 层。

其中 BatchedNestedDecoderpush_n_valids / push_n_invalids(L231-L253)同样采用"waiting 计数 + 交替时 flush"的批处理模式,与 Nested 自身的批量计数互为呼应。

columns_to_iter_recursive:外层递归驱动

nested.rscolumns_to_iter_recursive(L15-L245)是嵌套读取的驱动函数,递归方向与文档描述一致:

  • 非嵌套叶子(!field.dtype().is_nested()):直接走扁平入口 page_iter_to_arraysimple.rs);
  • List / LargeList:先 init.push(InitNested::List(field.is_nullable)),递归读元素列,再用 create_listnested 状态弹出 offsets/validity 组装 ListArray(L30-L44);
  • FixedSizeList:同上,InitNested::FixedSizeList(nullable, width)(L45-L59);
  • Struct:先递归读最后一个字段(因为列与类型列表按从后往前切分),弹出 struct 层的 lengthstruct_validity(L93-L96),再逆序读其余字段,debug_assert 校验各字段弹出的 struct validity 一致(L106-L113),最终组装 StructArray(L122-L134);
  • Map:按 List 语义递归后 create_map(L136-L150);
  • Dictionary(categorical/enum)与 Extension 类型也各自归入递归流程(L152-L238)。

init: Vec<InitNested> 参数沿递归层层追加,叶子处经 init_nested()nested_utils.rs)实例化为 NestedState,这正是文档所说"嵌套 parquet 字段初始被递归遍历以收集类型与可空性"的代码形态。

补充:扁平路径的 chunk 切分模型

deserialize 子模块的设计笔记(deserialize/README.md)进一步说明了扁平路径的迭代语义,可作为理解嵌套路径的基础:

  • 私有入口 simple::page_iter_to_arrays 接收一个已解压编码页的流式迭代器 Pages、Parquet 源类型、目标 Arrow DataTypechunk_size,返回 ArrayIter
  • 首次 pull 时构建 PageState,之后按 chunk_size 与页内行数的关系决定是继续拉页攒 chunk,还是把页内剩余行切入 FIFO 队列;
  • 页内数据布局为 [rep levels][def levels][non-null values],扁平列将页拆成"def levels 迭代器 + non-null values 迭代器"两条流,且"对非嵌套类型,def levels 就是与 Arrow 表示一致的位图,validity 直接 extend";
  • 嵌套类型则要求"行不能跨页切分"这一 Parquet 不变量(每个 page 必含完整行),因此 chunk_size 只驱动最外层 group 的行数,内层需要的项数由递归推算。

过滤(Filter)在嵌套路径上的现状

阅读 nested_utils.rscollect_nested 可以看到 Filter 各变体的处理差异:None 全量收集(num_collects = usize::MAX)、Range 按行区间跳过/收集、Mask 用顶层位图迭代驱动 skip/collect,而 Filter::Predicate(_) 目前是 todo!()(L622)。同样地,struct 分支在 nested.rs 断言 !matches!(&filter, Some(Filter::Predicate(_))),注释为 "This definitely does not support Filter predicate yet"。因此在当前仓库中,对嵌套列做谓词下推式读取尚未实现,仅 Range/Mask 过滤可用——使用嵌套类型谓词过滤时应以此为准。

工程细节与可观测性

几个源码中可见的工程决策值得注意:

  • 批量化位图扩展NestedBatchedNestedDecoder 都只对"连续 valid/invalid 段"计数、交替时才写位图。Nested::push 中的注释(nested_utils.rs)用 [{x:[1]}, None, {x:[1,2]}] 的例子说明了为何"没有 validity mask 的 list 层仍需为 null struct 插入 invalid 项"——Arrow 表示中 offsets 为 [0,1,1,3],即使 list 层无 validity 也要占位。
  • Primitive 层的零容量分配Nested::primitive 以 0 容量创建可选 BitmapBuilder(L33-L48),注释说明 primitive 的 validity 由 decoder 侧跟踪,此处仅用 Option 的有无标记可空性。
  • 解码指标utils/mod.rs 中的 PageDecoder::new 会读取环境变量 POLARS_PARQUET_METRICS=1OnceLock 惰性初始化),开启后逐页累计压缩/解压字节数与解压/解码微秒数,可用于诊断 Parquet 读取瓶颈。

小结

这份模块级设计文档用两条"Observations"概括了 Polars Parquet 读取器的核心取舍:

  1. 扁平可选列:max rep = 0 且 max def = 1 时,def levels 与 Arrow 位图 LSB 位序等价,validity 解码退化为 memcpy / extend_from_bitmap,配合 null_count == 0 的整页跳过,是读取热路径上最便宜的优化;
  2. 嵌套列:以 Nested/NestedState 为载体,把 Hybrid RLE 展开的 rep/def 整数流逐值翻译成 list offsets 与 struct/list validity,再通过 columns_to_iter_recursive 自内向外的递归 pop 重建 StructArray/ListArray;代价是放弃位图快路径并两遍消费 def levels。

对需要处理 List/Struct/Map 嵌套数据的 Polars 用户而言,这一实现意味着嵌套 Parquet 读取的正确性锚定在 Dremel levels 上;对阅读源码的贡献者而言,nested_utils.rsdecode_nestednested.rscolumns_to_iter_recursive 是理解该机制的两个必读入口。

登录后查看全文
热门项目推荐
相关项目推荐