Polars Parquet 读取器源码解析:嵌套类型反序列化与 Definition Levels 位图等价优化
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.rs 与 nested.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_validity,null_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,再递归弹出
文档给出的算法分两个阶段:
- 收集阶段:嵌套字段先被递归遍历,收集每一层"是 Struct 还是 List、required 还是 optional",存入
nested_info: Vec<Box<dyn Nested>>(在当前仓库中对应具体类型Nested,非 trait object)。Nested根据类型与可空性接收 def/rep levels,并据此累积每层的 offsets 与 validity。 - 弹出阶段:一个字段读完之后,从
nested_info中递归地pop,逐层构建StructArray或ListArray。
文档最后总结出嵌套路径相对扁平路径仅有的两点差异,这两点在源码中都有明确对应:
- 不再利用 bitmap 优化,需要把 repetition 与 definition levels 反序列化为整数值;
- 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_repeated、rep_delta = is_repeated,并预填初始值 0。这组阈值是后续逐值判断"本值在 depth 是否已定义"(def >= def_levels[depth])和"是否属于新行"(rep == 0)的依据。
collect_nested:把 rep/def 解码为整数并做两遍消费
文档说的"反序列化为 i32"与"def levels 消费两次"体现在 PageDecoder::collect_nested(nested_utils.rs):
- 每页调用
level_iters(L564-L576):从页数据中split_buffer拆出 def/rep 缓冲区,用列描述符中的max_def_level/max_rep_level计算 bit width,构建两个HybridRleDecoder; collect_level_values(L371-L391)把 Hybrid RLE 解码结果物化为Vec<u16>——这就是文档第 1 点:不再走"LSB 位图直接 memcpy"快路径,而是把 levels 完整展开为逐值整数;decode_nested(L409-L561)逐值消费:对每一对(def, rep),沿深度循环判断该层是否已定义、是否有效,调用nest.push(length, is_valid)累积各层的 offsets 与 validity 计数——这是 def levels 的第一次消费(扩展nested_info);同时通过BatchedNestedDecoder(L214-L284)生成叶子值的leaf_filter与leaf_validity两个位图——这是第二次消费(扩展值/可空性);- 最后
utils::State::new_nested+state.decode(L672-L683)用叶子 validity/filter 真正解码 primitive 值; - 解码完成后执行
nested_state.pop()弹出 Primitive 层(L689-L693),留下待上层columns_to_iter_recursive继续弹出 List/Struct 层。
其中 BatchedNestedDecoder 的 push_n_valids / push_n_invalids(L231-L253)同样采用"waiting 计数 + 交替时 flush"的批处理模式,与 Nested 自身的批量计数互为呼应。
columns_to_iter_recursive:外层递归驱动
nested.rs 的 columns_to_iter_recursive(L15-L245)是嵌套读取的驱动函数,递归方向与文档描述一致:
- 非嵌套叶子(
!field.dtype().is_nested()):直接走扁平入口page_iter_to_array(simple.rs); List/LargeList:先init.push(InitNested::List(field.is_nullable)),递归读元素列,再用create_list从nested状态弹出 offsets/validity 组装 ListArray(L30-L44);FixedSizeList:同上,InitNested::FixedSizeList(nullable, width)(L45-L59);Struct:先递归读最后一个字段(因为列与类型列表按从后往前切分),弹出 struct 层的length与struct_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 源类型、目标 ArrowDataType与chunk_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.rs 的 collect_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 过滤可用——使用嵌套类型谓词过滤时应以此为准。
工程细节与可观测性
几个源码中可见的工程决策值得注意:
- 批量化位图扩展:
Nested与BatchedNestedDecoder都只对"连续 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=1(OnceLock惰性初始化),开启后逐页累计压缩/解压字节数与解压/解码微秒数,可用于诊断 Parquet 读取瓶颈。
小结
这份模块级设计文档用两条"Observations"概括了 Polars Parquet 读取器的核心取舍:
- 扁平可选列:max rep = 0 且 max def = 1 时,def levels 与 Arrow 位图 LSB 位序等价,validity 解码退化为 memcpy /
extend_from_bitmap,配合null_count == 0的整页跳过,是读取热路径上最便宜的优化; - 嵌套列:以
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.rs 的 decode_nested 与 nested.rs 的 columns_to_iter_recursive 是理解该机制的两个必读入口。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00