首页
/ Polars compute 内核设计原则:polars-arrow 计算模块的 API 契约、不可变性与实现规范

Polars compute 内核设计原则:polars-arrow 计算模块的 API 契约、不可变性与实现规范

2026-09-05 09:31:20作者:裘旻烁

本篇技术指南基于 polars-arrow 计算模块的设计文档(crates/polars-arrow/src/compute/README.md),完整解读该模块面向分析场景的 13 条内核设计原则,并结合 aritybitwiseconcatenateutils 等真实源码逐条印证这些约定在实现中的落地方式。读完本文,你将理解 Polars 底层数组计算内核的接口契约(何时返回错误、如何传递所有权、如何做类型分派与特性门控),并能在阅读或扩展 polars-arrow 时准确把握每条约定背后的工程动机。

1. 模块定位:compute 提供分析场景的独立算子

polars-arrow/src/compute/ 是 polars-arrow 中存放"分析中常见的独立运算"的模块。其模块文档(mod.rs)开宗明义地给出了该模块的总体设计:每个算子都提供两个接口——静态类型版本和动态类型版本

  • 静态类型版本接收具体数组类型(如 PrimitiveArray<T>),编译期确定行为;
  • 动态类型版本接收 &dyn Array,在运行时检查逻辑类型,类型不支持时返回错误;
  • 部分动态类型算子还提供辅助函数 can_*,预先判断该算子能否作用于某个 DataType

当前模块包含以下子模块与文件(见 mod.rs):

文件 / 子模块 职责 特性门控
aggregate/ 聚合与内存估算(如 estimated_bytes_size compute_aggregate
arity.rs 一元/二元逐元素操作内核
arity_assign.rs 原地(in-place)一元/二元操作内核
bitwise.rs and / or / xor / not 及标量变体 compute_bitwise
concatenate.rs 多数组拼接为单一数组
decimal.rs 十进制数计算 dtype-decimal
temporal.rs 时间类型计算 compute_temporal
utils.rs 长度校验、validity 位图合并等工具函数

设计文档(README.md)声明该模块由"分析中常见的独立运算(independent operations common in analytics)"组成,并围绕这些运算确立了一组 MUST/SHOULD 级别的设计准则。下文按主题完整继承这些准则,并逐条对照源码。

2. 错误处理契约:何时必须返回错误

设计文档对错误处理提出了分层次的约定:

  1. MUST:当参数不正确,或运算会产生可预期错误(例如除以零)时,API 必须返回错误;
  2. MAY:当运算发生溢出时(如 i32 + i32),API 允许返回错误;
  3. SHOULD:当对可变大小容器的操作可能超过 usize 的最大值时,API 应当返回错误。

动态类型内核是这套契约最直接的体现。以 concatenate 为例(concatenate.rs):

/// Concatenate multiple [`Array`] of the same type into a single [`Array`].
pub fn concatenate(arrays: &[&dyn Array]) -> PolarsResult<Box<dyn Array>> {
    if arrays.is_empty() {
        polars_bail!(InvalidOperation: "concat requires input of at least one array")
    }
    if arrays.iter().any(|array| array.dtype() != arrays[0].dtype()) {
        polars_bail!(InvalidOperation: "It is not possible to concatenate arrays of different data types.")
    }
    concatenate_unchecked(arrays)
}

这里完整演示了三类"参数不正确"的处理:空输入、数据类型不一致都会通过 polars_bail! 返回 InvalidOperation 错误,而不 panic。这也是文档第 12 条约定(动态类型数组必须以 &dyn Array 传入)与错误契约的结合点——concatenate 正是接收 &[&dyn Array] 的动态类型算子,其内部按 PhysicalType 分发到各具体类型的拼接实现(concatenate.rs),对尚未实现的 UnionMapDictionary 类型直接 unimplemented!()

静态类型内核则采用"契约上界 + 快速路径"的写法。arity::binary 的文档注释明确其错误条件(arity.rs):

/// Applies a binary operations to two primitive arrays.
///
/// # Errors
/// This function errors iff the arrays have a different length.
///
/// # Implementation
/// This will apply the function for all values, including those on null slots.
/// This implies that the operation must be infallible for any value of the
/// corresponding type.
pub fn binary<T, D, F>(
    lhs: &PrimitiveArray<T>,
    rhs: &PrimitiveArray<D>,
    dtype: ArrowDataType,
    op: F,
) -> PrimitiveArray<T>

其实现调用 check_same_len(lhs, rhs).unwrap()arity.rs L57),而 check_same_len 才是真正产出错误的地方(utils.rs L112-L118):

// Errors iff the two arrays have a different length.
#[inline]
pub fn check_same_len(lhs: &dyn Array, rhs: &dyn Array) -> PolarsResult<()> {
    polars_ensure!(lhs.len() == rhs.len(), ComputeError:
            "arrays must have the same length"
    );
    Ok(())
}

可以推断出这里的设计权衡:错误判定逻辑(check_same_len)统一返回 PolarsResult,满足 MUST 级契约;而逐元素内核在已确认调用方受控的场景下选择 unwrap() 快速路径。同时注意文档注释强调"对 null 槽位也执行操作"——这解释了为什么要求传入闭包必须对任意取值都不失败(infallible),否则内核会 panic。

关于溢出 MAY 级约定:arity::binary 的闭包签名 F: Fn(T, D) -> T 将溢出语义完全交给闭包实现者选择(可用 +,也可用 checked_add 后由上层转换为错误),模块本身不做强制,这与"MAY error when an operation overflows"的弹性表述一致。

3. 不可变性与所有权:kernels MUST NOT 产生副作用、MUST NOT 接管参数

文档中最强硬的两条约定是:

  • kernels MUST NOT have side-effects(不得有副作用);
  • kernels MUST NOT take ownership of any of its arguments(不得接管任何参数的所有权,一切参数必须按引用传递)。

arity::unary / arity::binary 是这两条约定的标准样板:参数均为 &PrimitiveArray<_>,返回全新构造的 PrimitiveArray<O>,对输入零修改(arity.rs)。validity 位图也遵循"引用 + clone"模式——array.validity().cloned() 只是克隆 Option<Bitmap> 的内部缓冲共享,而非转移所有权。

如果确实需要就地写入,模块提供的是显式命名、语义明确的另一族 APIarity_assign.rsunary 接收 &mut PrimitiveArray<I>,并通过 COW(写时复制)语义决定是否真正原地修改(arity_assign.rs L20-L33):

pub fn unary<I, F>(array: &mut PrimitiveArray<I>, op: F)
where
    I: NativeType,
    F: Fn(I) -> I,
{
    if let Some(values) = array.get_mut_values() {
        // mutate in place
        values.iter_mut().for_each(|l| *l = op(*l));
    } else {
        // alloc and write to new region
        let values = array.values().iter().map(|l| op(*l)).collect::<Vec<_>>();
        array.set_values(values.into());
    }
}

只有当底层缓冲独占(非共享)时才原地写,否则分配新缓冲。其二元版本 binary 对 validity 与 values 都做同样的 COW 分支处理,源码注释里还记录了一个实测结论(arity_assign.rs L58-L60):"先 memcpy 再原地赋值"比直接分配新缓冲慢约 2 倍,因此共享时选择整体分配新区域。这说明"无副作用"约定并非一刀切禁止就地写,而是要求:就地写必须有独立命名的 API 族、有清晰的 &mut 边界,且行为对缓冲共享性透明。

4. 类型系统约定:用逻辑类型裁决适用性,动态类型走 &dyn Array

文档中两条关于类型的约定:

  • SHOULD use the arrays' logical type:内核应依据数组的逻辑类型判断自身能否被应用。文档给出的例子是 Date32 + Date32 没有意义,因此不应该被实现;
  • 动态类型内核 MUST 以 &dyn Array 接收数组

第一条在 concatenate 的分发表中得到体现:它按 dtype().to_physical_type() 匹配出 NullBooleanPrimitive(_)Utf8ListStruct 等分支(concatenate.rs),对每个物理类型选择唯一的实现,而不是对任意类型强行泛化——这正是"依据逻辑类型裁决适用性"的工程含义:不支持的组合就明确拒绝(返回错误或 unimplemented!),而不是产出一个语义错误的结果。

第二条即"当 API 返回 &dyn Array 时必须返回 Box<dyn Array>"这一约定(详见第 7 节),而 concatenate 的返回值正是 PolarsResult<Box<dyn Array>>

此外,validity 合并工具展示了另一层类型层面的纪律:combine_validities_and / combine_validities_or / combine_validities_and3 / combine_validities_and_many 全部以 Option<&Bitmap> 为输入(utils.rs),逐一穷举 None/Some 组合,保证任何含 null 的操作都能正确传播有效性。combine_validities_and_many 甚至针对 0/1/2/3 个有效位图做了特化分支,多于一三个时改用 fast_iter_u64 按 64 位块求与(utils.rs L54-L109)。

5. 位读取纪律与自动向量化:只允许通过 Bitmap 读位、优先 from_trusted_len_iter

文档规定:

  • Kernels MUST NOT use any API to read bits other than the ones provided by Bitmap
  • Implementations SHOULD aim for auto-vectorization, which is usually accomplished via from_trusted_len_iter

前者在 utils.rs 中得到贯彻:所有 validity 运算(bitandbitorternaryand_notpush_bitchunk)都来自 crate::bitmap 模块提供的 API(utils.rs L7),compute 模块自身没有任何绕过 Bitmap 直接操作位的手工代码。这样位运算的内存对齐、越界与尾位处理都收敛在 polars-arrow/src/bitmap/ 一处维护。

后者的意图是:逐元素映射优先用 iter().map(...).collect::<Vec<_>>() 这类可被 Rust 自动向量化(auto-vectorization)的写法,长度可信的迭代器场景则走 from_trusted_len_iterarity::binary 的取值路径就是这种风格(arity.rs L61-L67):

let values = lhs
    .values()
    .iter()
    .zip(rhs.values().iter())
    .map(|(l, r)| op(*l, r))
    .collect::<Vec<_>>()
    .into();

简洁的 zip + map + collect 形态使编译器有机会对 T 的逐元素运算生成 SIMD 代码,这正是文档"auto-vectorization"约定的典型实现形态。

6. 缓冲操作偏好:通过 clone / slice / iterator 操作 Buffer 与 Bitmap

文档约定 Kernels SHOULD be implemented via clonesliceBufferBitmapVecMutableBitmap 提供的 iterator API。这条约定的目的是让内核尽量复用已有的缓冲共享与零拷贝机制,而不是自行拷贝数据。

源码中有大量对应实例:

  • concatenate_primitive 直接用 out.extend_from_slice(arr.values()) 拼接值缓冲,再以 Buffer::from(out) 一次性构建结果(concatenate.rs L133-L144);
  • concatenate_view 在多个 view 数组共享同一缓冲集时,直接 Buffer::clone(first_arr.data_buffers()) 复用缓冲、只复制 views,避免整段数据拷贝(concatenate.rs L220-L242);
  • arity_assign::binary 对 validity 的合并则区分 Either::Left(immutable)(共享、需新分配)与 Either::Right(mutable)(独占、原地 & 运算)两种情况(arity_assign.rs L65-L76)。

7. 返回值约定:返回 &dyn Array 时必须返回 Box<dyn Array>

文档最后一条约定及其理由值得单独强调:

When an API returns &dyn Array, it MUST return Box<dyn Array>. The rationale is that a Box is mutable, while an Arc is not. As such, Box offers the most flexible API to consumers and the compiler. Users can cast a Box into Arc via .into().

核心逻辑是:Box<dyn Array> 既可转 Arc(通过 .into())也可继续可变操作,而一旦 API 返回 Arc,调用方就永远失去了"独占/可变"的可能性。因此返回类型选择 Box信息量最大、对消费者最灵活的契约。concatenate 的签名 -> PolarsResult<Box<dyn Array>>concatenate.rs L16)严格遵循了这一约定,其内部各分支通过 to_boxed() 统一收口。

8. 特性门控:依赖外部依赖的实现必须 feature-gate

文档要求 Implementations MUST feature-gate any implementation that requires external dependenciesmod.rs 中的模块声明完整展示了这一做法:

#[cfg(feature = "compute_aggregate")]
#[cfg_attr(docsrs, doc(cfg(feature = "compute_aggregate")))]
pub mod aggregate;
...
#[cfg(feature = "compute_bitwise")]
#[cfg_attr(docsrs, doc(cfg(feature = "compute_bitwise"))]
pub mod bitwise;
...
#[cfg(feature = "compute_temporal")]
pub mod temporal;

对应的开关定义在 crates/polars-arrow/Cargo.toml(第 103–114 行附近):compute_aggregate = []compute_bitwise = []compute_temporal = [],且聚合进一个默认特性集合中。每个门控模块同时用 cfg_attr(docsrs, doc(cfg(...))) 在文档站点标注其依赖的特性,方便使用者按图索骥。decimal 模块则门控在 dtype-decimal 下。

9. 动态类型算子的实例:estimated_bytes_size

aggregate/memory.rs 中的 estimated_bytes_size 是"动态类型接口 + 类型分派 + 文档级实现说明"的又一范例(memory.rs L33-L44):

/// Returns the total (heap) allocated size of the array in bytes.
/// ...
pub fn estimated_bytes_size(array: &dyn Array) -> usize {
    use PhysicalType::*;
    match array.dtype().to_physical_type() {
        Null => 0,
        Boolean => { ... },
        Primitive(primitive) => with_match_primitive_type_full!(primitive, |$T| { ... }),
        ...
    }
}

它的文档注释本身就是一份实现事实:多个数组可能共享缓冲,因此两个数组的大小之和不等于各自估算值的和(StructArray 的大小是上界);数组被 slice 后,由于缓冲不变,该函数返回的是"可见大小"而非总容量。这些限定恰好呼应了设计文档第 3 条精神——可预期的边界行为必须在文档中写明,而不是留给调用方猜测。该函数被上层用于判断内存占用(例如触发溢写/去重的阈值计算),是 compute 模块服务分析场景的典型入口之一。

10. 实践小结:按设计文档阅读与扩展 compute 模块

综合设计文档与源码证据,可以提炼出面向开发者的三条实操结论:

  1. 调用侧:优先使用动态类型入口(接收 &dyn Array、返回 PolarsResult<Box<dyn Array>> 的算子,如 concatenate);需要判断类型适用性时先查是否存在对应的 can_* 辅助函数(见 mod.rs 的模块说明)。静态类型入口(aritybitwise)返回裸数组、错误契约以文档注释声明,适合编译期类型确定、追求最快路径的场景。
  2. 实现侧:新内核应遵循"参数全引用、无副作用、null 槽位也参与计算且闭包必须 infallible、位运算只走 Bitmap API、缓冲操作走 clone/slice/iterator"五条铁律;需要溢出/长度校验时用 check_same_lenpolars_ensure! 返回 PolarsResult,而不是散落各处的手工判断。
  3. 扩展侧:若新实现引入外部依赖,必须像 bitwise/temporal/aggregate 一样在 Cargo.toml 中定义独立 feature 并在 mod.rs#[cfg] 门控;若新算子对某些逻辑类型无意义(如日期加日期),按约定直接不实现,让动态类型入口返回错误,而不是硬凑一个能编译的结果。

这套"契约先行(MUST/SHOULD 文档化)+ 双接口(静态/动态类型)+ COW 原地写单独命名 + 位操作收敛于 Bitmap"的设计,是 polars-arrow compute 模块保持各算子独立、可组合、可被上层引擎(如 polars-core 的 ChunkedArray 运算与 polars-ops 的算子实现)安全复用的基础。

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

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.79 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
988
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384