首页
/ Polars 之 polars-arrow IO 模块:基于格式分目录、io_* 特性门控与元数据/数据分离的读写架构设计

Polars 之 polars-arrow IO 模块:基于格式分目录、io_* 特性门控与元数据/数据分离的读写架构设计

2026-09-05 09:29:23作者:蔡怀权

本文以 IO 模块设计文档 为核心,系统讲解 polars-arrow 内部 io 模块的整体设计规则:按格式分目录组织代码、用 io_* 前缀特性门控外部依赖、强制 re-export 外部 API 的“Cargo.toml 即够”原则、read/write 子目录约定,以及“元数据与数据分离”“IO 受限与 CPU 受限操作分离”两条关键实现原则,并结合 ipcavro 两个真实格式模块的源码印证这些规则是如何落地的。读完本文,你将能够理解 polars 底层 Arrow 读写层的模块组织逻辑、各特性开关的依赖关系,以及为何 FileReader/FileWriterStreamReader/StreamWriter 这样拆分为独立公共 API。

一、设计文档定下的五条模块规则

polars-arrow 仓库在 IO 模块文档 中以“Rules”一节明确了整个 io 模块必须遵循的架构约束,这是后文所有源码分析的依据:

  1. 每个目录对应一个具体格式:如 csvjson。当前仓库中该模块下实际存在 ipc(Arrow IPC 格式)与 avro(Apache Avro 格式)两个格式目录,另有跨格式公共件 iterator.rswrite_owned.rs
  2. 依赖外部 crate 的目录必须特性门控,特性名以 io_ 为前缀:这保证了核心 crate 在默认构建时不引入任何外部格式依赖;
  3. 模块必须把所用外部依赖的公共 API 一并 re-export:例如某个模块的 API 签名里出现 writer: &mut csv::Writer<W>,则该模块必须包含 pub use csv::Writer;。文档给出的理由是“把该 crate 加入 cargo.toml 之后就应该可以直接使用它”,调用方不需要再额外添加同一个外部依赖来匹配类型;
  4. 每个格式目录应包含 readwrite 两个子目录,分别承载“从该格式读入”与“写出到该格式”的功能;
  5. 基础模块应声明 pub use read;pub use write;,使 io::ipc::readio::ipc::write 形成稳定的公共入口。

此外,文档还提出两条面向“实现质量”的 SHOULD 级原则,分别对应下两节的源码印证:

  • 数据读取与元数据读取分离:schema 读取/推断应是独立函数;读取“数据”的函数应接收一个(通常预先读好的)schema;
  • IO 受限与 CPU 受限操作分离:应提供“消费 Read 实现者并输出 raw 结构(如仍压缩/序列化状态的数据)”的函数,再提供“消费 raw 结构并转换为 Arrow 数组”的函数,且两者都作为独立公共 API 暴露,让调用方自行权衡 CPU 与 IO 的平衡。

二、特性门控:io_* 前缀特性与依赖链

模块入口 io/mod.rs 只有十余行,但完整体现了规则 2:

#[cfg(feature = "io_ipc")]
pub mod ipc;

#[cfg(feature = "io_avro")]
pub mod avro;

pub mod iterator;
pub mod write_owned;

格式模块整体被 #[cfg(feature = "...")] 包裹,默认不编译;而 iteratorwrite_owned 这类不依赖外部格式的通用工具始终编译。

再看 polars-arrow/Cargo.toml 中 features 一节,每个 io_* 特性都精确地把“可选依赖 + 对应 error 特性”串起来:

特性 依赖内容 说明
io_ipc arrow-formatpolars-error/arrow-format Arrow IPC(FlatBuffers 协议)基础读写
io_ipc_compression lz4zstdio_ipc 在 IPC 之上叠加 LZ4/Zstd 压缩
io_flight io_ipcarrow-format/flight-dataasync-streamfuturestokio 异步 Arrow Flight 转换
io_avro avro-schemapolars-error/avro-schema Avro 读写
io_avro_compression avro-schema/compression Avro 压缩支持
io_avro_async avro-schema/async Avro 异步读取
full io_ipcio_flightio_ipc_compressionio_avroio_avro_compressionio_avro_asynccomputeserdechrono-tz 一键开启全部 IO 能力(docs.rs 构建即使用 full

这种“基础特性 + 压缩特性 + 异步特性”的分层组织,让下游 crate 可以只编译自己需要的格式路径——这正是规则 2 想要的效果:依赖外部 crate 的目录被隔离,编译产物与功能范围由 feature 决定。

三、read / write 子目录与 re-export 规则落地

3.1 目录结构完全遵循约定

ipcavro 两个格式目录均按规则 4 组织为 read + write 子目录:

  • io/ipc/read/(含按 Arrow 类型拆分的逐类型反序列化,read/array/ 下覆盖 primitive、utf8、binary、binview、list、dictionary、union、map 等全部物理类型)、write/serialize/ 下对称地按类型拆分序列化器)、write2/(新的写入路径:arrayfootermessageschema)与 append/
  • io/avro/read/deserializenestedschemautil)与 write/schemaserialize)。

avro/mod.rs 是规则 5 的直接体现:

pub use avro_schema;   // re-export 外部依赖,规则 3
pub mod read;
pub mod write;

其中 pub use avro_schema; 恰好是文档中“必须 re-export 外部依赖 API”这条规则的现实案例:avro_schema::file::FileMetadataavro_schema::read::BlockStreamingIterator 等外部类型都会出现在本模块公共 API 的签名中,若不 re-export,调用方就必须自己依赖 avro-schema crate 并保证版本一致。IPC 侧同理,io/ipc/mod.rspub use arrow_format as format; 把外部 crate arrow-formatformat 别名公开,io/ipc/write/mod.rs 进一步 pub use arrow_format::ipc::{Block, KeyValue, KeyValueRef}; 并批量 re-export common 中的 CompressionWriteOptionsencode_record_batch 等写入核心 API。

四、“元数据与数据分离”:先读 footer/inferrer,再按 schema 读数据

文档要求“schema 读取或推断应是独立函数”“读数据的函数应消费一个预先读好的 schema”。两个格式模块都严格执行了这一点:

IPC 文件io/ipc/read/file.rs 定义了 FileMetadata 结构体,聚合了从文件 footer 读出的全部元信息:

pub struct FileMetadata {
    pub schema: ArrowSchemaRef,                    // footer 中的 schema
    pub custom_schema_metadata: Option<Arc<Metadata>>,
    pub custom_metadata: Option<Arc<Metadata>>,
    pub ipc_schema: IpcSchema,                      // IpcField + 端序
    pub blocks: Vec<arrow_format::ipc::Block>,      // 各记录块的文件区段
    pub dictionaries: Option<Vec<arrow_format::ipc::Block>>,
    pub size: u64,
}

配套的公开函数 read_file_metadata<R: Read + Seek>(reader: &mut R) -> PolarsResult<FileMetadata> 只负责读元数据;io/ipc/read/reader.rs 中的 FileReader::new(reader, metadata, projection, limit)接收这个预读好的 FileMetadata 作为构造参数——“数据读取消费预读 schema”的要求在函数签名层面就完成了。同文件的 get_row_count 也展示了同一思路:先 read_footer_lenread_footerdeserialize_footer_blocks,再按 blocks 求和行数,全程不触碰数据块本身。

Avroio/avro/read/mod.rsReader::new 的签名同样要求外部先完成“元数据阶段”:

pub fn new(
    reader: R,
    metadata: FileMetadata,        // avro_schema 读出的文件元数据
    fields: ArrowSchema,           // 由 infer_schema 推断出的 Arrow schema
    projection: Option<Vec<bool>>,
) -> Self

其中 fields 来自独立的推断函数 io/avro/read/schema.rs 中的 pub fn infer_schema(record: &Record) -> PolarsResult<ArrowSchema>——这正是文档中“schema 推断应是单独函数”的示例。写入侧也保持对称:io/avro/write/mod.rs 把“schema 转 Avro Record”(to_record)、“按 schema 构造序列化器”(new_serializer/can_serialize)与“把一组列式 BoxSerializer 转置为行式 Block 字节”(serialize)拆成独立函数,调用者可以决定哪些步骤自己做。

五、“IO 受限与 CPU 受限分离”:raw 层与 Arrow 层的两级 API

文档对这一原则的表述是:提供“消费 Read 实现者、输出 raw 结构(如仍压缩/序列化)”的函数,再提供“消费 raw 结构并转换为 Arrow”的函数,两者都作为独立公共 API 开放。

IPC 读路径io/ipc/read/common.rs 中体现了这一分层。以文件读取为例,调用链大致为:

  1. get_message_from_block(reader, block, &mut scratch) —— 纯 IO 阶段:按 footer 给出的 Block 定位(offset/meta_data_length)并读出一个 raw 的 FlatBuffers MessageRef,输出仍是“序列化/布局好的字节”,尚未产生任何 Arrow 数组;
  2. read_dictionary_block(...) / read_batch(...) —— CPU 阶段:消费上一步的 raw message,按 schema 与 dictionary 解码出 RecordBatch

reader.rsnext_record_batch 只做到第 1 层(返回 arrow_format::ipc::RecordBatchRef,即 raw 视图),而 Iterator for FileReadernext 则走完整解码路径得到 RecordBatchT<Box<dyn Array>>。两个层级都对外可用,调用方既能“只读消息头”(例如 BlockReader 只需 record_batch_num_rows 就能拿到行数),也能做完整转换。此外 FileReader 暴露了 take_scratches/set_scratches,允许把解码过程中复用的 data_scratch/message_scratch 缓冲移交给下一个读写器,把 CPU 侧的分配成本也交给调用方控制——这正是“让消费者自行平衡 CPU-bounds 与 IO-bounds”的延伸。

Avro 读路径同样分层:block_iterator(reader, compression, marker)(来自 re-export 的 avro_schema)负责 IO——把文件切分成带压缩信息的 raw Block 流;io/avro/read/deserialize.rspub fn deserialize 负责 CPU——消费一个 block 与预推断的 ArrowSchema/avro_fields 产出 RecordBatchTReader 迭代器只是把这两步粘合起来,两步各自独立可调用。

Avro 写路径则展示了 raw 层的完整定义:io/avro/write/serialize.rs 中每个 BoxSerializernext() 产出一行数据的原始字节,serialize 将其转置进 Block——Block 本身就是文档所指的 raw 结构(含 datanumber_of_rows,后续交给 avro-schema 的压缩/落盘),而把列式 Arrow 数组“切成逐行序列化器”是另一层职责,调用方可以自由选择在哪里结束。

六、两种 IPC 读写器:File 变体与 Stream 变体

io/ipc/mod.rs 的模块文档说明了 IPC 模块的整体形态:底层用 FlatBuffers 作为二进制协议(通过 re-export 的 arrow-format 处理),并定义了格式魔数常量(ARROW_MAGIC_V1 = b"FEA1"ARROW_MAGIC_V2 = b"ARROW1"、带 padding 的 8 字节版本及 CONTINUATION_MARKER),供读取端识别与校验文件头。

文档进一步指出该模块的读写 API 分两个变体:

  • FileReader <-> FileWriterwrite/writer.rs):包装同时实现 Read + SeekT,因此 File 变体支持随机访问——读取时先读 footer 得到 blocks,可以按块跳读、skip_blocks_till_limit 跳过前 N 个 block,配合 projection 做列裁剪与 limit 限行;
  • StreamReader <-> StreamWriterwrite/stream.rsread/stream.rs):只需 Read(流式),按写入顺序“先进先出”顺序消费,适合管道/网络流场景。

IPC 协议本身要求预定义 schema(只传 record batch 或 dictionary batch),文档也解释了这带来的收益:序列化数据的结构恒定,反序列化时字节已“摆好位置”,速率显著提升。FileReader::new 还支持 projection: Option<Vec<usize>>(要求索引递增)与 limitunchecked() 则提供关闭一致性检查的快速路径(UnsafeBool),供可信数据源跳过昂贵校验——这些都属于模块在“平衡 CPU/IO 成本”上给调用方留出的旋钮。

写入侧,io/ipc/write/mod.rsdefault_ipc_field/default_ipc_fields 负责递归遍历 schema、为每个 Dictionary 类型字段分配递增 dictionary_id 并生成 IpcField 树(List/Map/Struct/Union 递归子节点,叶子类型不分配 id),随后 encode_record_batchcommit_encoded_arraysdictionaries_to_encode 等函数完成“Arrow 数组 -> FlatBuffers 消息”的编码,全部作为 pub use 暴露给上层 crate(如 polars-io 的 IPC 读写层)。

七、小结:设计文档与仓库实现的一致性

回看 IO 模块文档 的五条规则与两条实现原则,可以在当前仓库中找到一一对应的实现证据:

设计文档要求 仓库中的落点
每个目录对应一个格式 io/ipc/io/avro/
外部依赖目录 io_* 特性门控 io/mod.rs#[cfg(feature = ...)]Cargo.toml 的 features 表
re-export 外部 API(“Cargo.toml 即够”) avro/mod.rspub use avro_schema;ipc/mod.rspub use arrow_format as format;
每目录含 read/write 子目录 ipc 与 avro 均为 read/ + write/ 结构
元数据/数据分离(独立 schema 推断函数) read_file_metadata + FileReader::new(reader, metadata, ...)infer_schema + avro::read::Reader::new(reader, metadata, fields, ...)
IO 与 CPU 操作分离、各自独立 API get_message_from_block(raw 消息)与 read_batch(Arrow 解码);block_iterator(raw Block 流)与 deserialize(Arrow 解码);BoxSerializer/serialize(raw Block 写入)

这套设计让 polars-arrow 的 io 模块保持了清晰的边界:新增一种格式只需新建一个目录、一个 io_<format> 特性,并沿用 read/write 子目录与两级(raw/Arrow)API 约定;下游 crate 则始终可以“加一行依赖就能用”。理解这一约定,是阅读或扩展 polars 中任何 IO 相关代码(包括上层 polars-io 的 IPC/Avro 封装)的前置基础。

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