Polars Expression Plugins 完全指南:用 Rust 编写原生速度的 DataFrame 表达式插件
【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars
本指南以 Polars 官方用户文档《Expression Plugins》为核心脉络,系统讲解如何在 Python 生态中把一段 Rust 函数编译为动态链接库,并注册成与 Polars 内置表达式几乎同速的表达式插件(Expression Plugin)。你将掌握从零搭建插件工程、编写#[polars_expr]自定义表达式、接收 kwargs、推导输出类型,以及插件在 Polars 引擎中被动态加载与调用的底层原理,并可直接复用仓库中完整的可运行示例工程。
什么是 Expression Plugins:在 Python 中运行近乎原生的 Rust 表达式
Polars 为「用户自定义函数」(UDF)场景提供了多种方案,而表达式插件(Expression Plugins)是官方推荐的首选方式。它的思路是:你在 Rust 侧实现一个普通函数,借助pyo3-polars提供的#[polars_expr]派生宏把它编译为一个导出符号的cdylib动态库,再通过 Python 侧的register_plugin_function把该函数以表达式形式注册进 Polars。当查询执行时,Polars 引擎会在运行时通过libloading动态链接你的函数,因此表达式运行速度几乎等同于原生表达式。
与常见的逐行回调型 UDF 不同,插件执行路径完全不经过 Python 解释器,因此不存在 GIL(全局解释器锁)争用问题——这既是性能的关键,也是插件可以无缝利用多线程的前提。一个已注册的插件表达式会继承 Polars 默认表达式的三大优势:
- 优化(Optimization):插件作为
Expr树中的一个节点参与逻辑计划与物理计划的优化(例如与谓词下推、投影下推协同); - 并行(Parallelism):声明为 elementwise 的插件可以在流式 / 并行执行引擎中被切分到多线程批量执行(见 crates/polars-stream 中的执行模型);
- Rust 原生性能(Rust native performance):核心计算保持在 Rust 侧以零开销方式运行。
从源码结构看,插件在计划层由FunctionExpr中的 plugin 分支表示,真正的运行时调用位于 crates/polars-plan/src/plans/aexpr/function_expr/plugin.rs,我们将在文末的「底层原理」一节还原其完整调用链。
第一个自定义表达式:Pig Latin 转换器
为了体验完整流程,官方文档以 Pig Latin 为例。Pig Latin 是一种把单词首字母移到末尾并追加ay的「人造语言」,例如pig会变成igpay。它足够简单,可以让你把注意力放在插件工程的搭建与注册机制上。
实际上,这个功能用纯 Polars 表达式也能实现:
pl.col("name").str.slice(1) + pl.col("name").str.slice(0, 1) + "ay"但一个专门的 Rust 函数会比这种字符串拼接表达式性能更好,而且它是学习插件机制的最佳入门案例。下面我们开始搭建插件工程。
搭建插件工程:Cargo.toml 与目录结构
首先创建一个新的 Rust 库,其Cargo.toml如下:
[package] name = "expression_lib" version = "0.1.0" edition = "2021" [lib] name = "expression_lib" crate-type = ["cdylib"] [dependencies] polars = { version = "*" } pyo3 = { version = "*", features = ["extension-module", "abi3-py310"] } pyo3-polars = { version = "*", features = ["derive"] } serde = { version = "*", features = ["derive"] }关键点说明:
crate-type = ["cdylib"]:必须将 crate 编译为 C 动态库,Polars 引擎才能以 C ABI 方式加载你导出的符号;pyo3-polars需开启derivefeature,它提供#[polars_expr]派生宏;serde用于自定义 kwargs 结构体的反序列化(见后文「接受 kwargs」章节);- 生产工程中建议把
polars/pyo3-polars/arrow等依赖收敛到工作区统一版本管理,仓库内示例 pyo3-polars/example/derive_expression/expression_lib/Cargo.toml 即采用workspace = true的方式引用同一版本库的polars、arrow、pyo3、rayon等依赖,避免与宿主 Polars 版本产生 ABI 不兼容。
库的顶层入口需要安装 Polars 的内存分配器,以保证数据在跨 FFI 边界传递时内存分配与释放策略一致:
// src/lib.rs use pyo3_polars::PolarsAllocator; mod distances; mod expressions; #[global_allocator] static ALLOC: PolarsAllocator = PolarsAllocator::new();PolarsAllocator由pyo3-polars提供(见 pyo3-polars/pyo3-polars/src/alloc.rs 及 src/lib.rs),在插件场景下把全局分配器设为 Polars 的分配器是官方示例工程的标准做法。
编写 Rust 侧的表达式函数
在src/expressions.rs中,先写一个把&str转换为 pig latin 的纯函数,再写一个暴露为表达式的包装函数。暴露函数必须添加#[polars_expr(output_type=DataType)]属性,且第一个参数必须是inputs: &[Series],返回PolarsResult<Series>:
// src/expressions.rs use polars::prelude::*; use pyo3_polars::derive::polars_expr; use std::fmt::Write; fn pig_latin_str(value: &str, output: &mut String) { if let Some(first_char) = value.chars().next() { write!(output, "{}{}ay", &value[1..], first_char).unwrap() } } #[polars_expr(output_type=String)] fn pig_latinnify(inputs: &[Series]) -> PolarsResult<Series> { let ca = inputs[0].str()?; let out: StringChunked = ca.apply_into_string_amortized(pig_latin_str); Ok(out.into_series()) }需要说明的若干实现细节:
inputs[0].str()?:从第一个输入列取出字符串类型的StringChunked,类型不符会直接传播PolarsResult错误;- 使用
apply_into_string_amortized而非apply_values:前者会复用同一个输出String缓冲区,避免为每行都分配新字符串,是处理字符串型 elementwise 转换的推荐 API; - 多输入、逐元素操作且输出为
String的场景,可以改用polars::prelude::arity中的binary_elementwise_into_string_amortized工具函数,它位于 crates/polars-core/src/chunked_array/ops/arity.rs。
Python 侧注册:文件夹命名与函数名必须匹配
Rust 侧到此就完成了。Python 侧需要建立与Cargo.toml中[lib] name(即expression_lib)同名的 Python 包目录,并在其中提供__init__.py。最终目录结构如下:
├── 📁 expression_lib/ # 名称必须与 Cargo.toml 的 "lib.name" 一致 │ └── __init__.py │ ├── 📁 src/ │ ├── lib.rs │ └── expressions.rs │ ├── Cargo.toml └── pyproject.toml这个同名目录正是插件动态库最终被maturin安装的位置。接着在__init__.py中注册新表达式:
# expression_lib/__init__.py from pathlib import Path from typing import TYPE_CHECKING import polars as pl from polars.plugins import register_plugin_function from polars._typing import IntoExpr PLUGIN_PATH = Path(__file__).parent def pig_latinnify(expr: IntoExpr) -> pl.Expr: """Pig-latinnify expression.""" return register_plugin_function( plugin_path=PLUGIN_PATH, function_name="pig_latinnify", args=expr, is_elementwise=True, )其中几个参数必须理解到位:
function_name必须与 Rust 侧的#[polars_expr]函数名完全一致(此处均为pig_latinnify),否则主 Polars 包无法解析到对应符号;plugin_path指向插件包所在目录(运行时会被解析为其中.so/.dll/.pyd动态库,解析逻辑见 py-polars/src/polars/plugins.py 中的_resolve_plugin_path与_is_dynamic_lib);is_elementwise=True告知 Polars 该函数是逐元素操作,从而允许引擎对它做批量切分与并行执行;而类似排序(sort)、切片(slice)这类改变数据整体形态的操作则不能声明为 elementwise。
编译与使用
在当前环境中安装maturin,然后编译并安装到当前虚拟环境:
pip install maturin maturin develop --release一切就绪后,表达式即可像内置表达式一样被使用:
import polars as pl from expression_lib import pig_latinnify df = pl.DataFrame( { "convert": ["pig", "latin", "is", "silly"], } ) out = df.with_columns(pig_latin=pig_latinnify("convert"))进阶:把插件挂载为自定义命名空间
除了「函数式调用」外,还可以通过 Polars 的命名空间注册 API 创建自定义命名空间(例如Expr.language),让用户以链式调用的风格编写:
out = df.with_columns( pig_latin=pl.col("convert").language.pig_latinnify(), )命名空间注册入口为polars.api.register_expr_namespace(及对应的register_series_namespace),其实现位于 py-polars/src/polars/api.py。官方示例工程 pyo3-polars/example/derive_expression/expression_lib/expression_lib/language.py 与 extension.py 中展示了如何把pig_latinnify、append_args等插件封装进命名空间类并注册,读者可直接对照阅读。
接受 kwargs:把普通参数传入插件函数
很多真实场景需要向插件传递普通(非 Series)参数。做法是定义一个Ruststruct,让它派生serde::Deserialize,并将该类型作为插件函数的第二个参数接收:
/// Provide your own kwargs struct with the proper schema and accept that type /// in your plugin expression. #[derive(Deserialize)] pub struct MyKwargs { float_arg: f64, integer_arg: i64, string_arg: String, boolean_arg: bool, } /// If you want to accept `kwargs`. You define a `kwargs` argument /// on the second position in you plugin. You can provide any custom struct that is deserializable /// with the pickle protocol (on the Rust side). #[polars_expr(output_type=String)] fn append_kwargs(input: &[Series], kwargs: MyKwargs) -> PolarsResult<Series> { let input = &input[0]; let input = input.cast(&DataType::String)?; let ca = input.str().unwrap(); Ok(ca .apply_into_string_amortized(|val, buf| { write!( buf, "{}-{}-{}-{}-{}", val, kwargs.float_arg, kwargs.integer_arg, kwargs.string_arg, kwargs.boolean_arg ) .unwrap() }) .into_series()) }Python 侧在注册时把同名 kwargs 一并传入:
def append_args( expr: IntoExpr, float_arg: float, integer_arg: int, string_arg: str, boolean_arg: bool, ) -> pl.Expr: """ This example shows how arguments other than `Series` can be used. """ return register_plugin_function( plugin_path=PLUGIN_PATH, function_name="append_kwargs", args=expr, kwargs={ "float_arg": float_arg, "integer_arg": integer_arg, "string_arg": string_arg, "boolean_arg": boolean_arg, }, is_elementwise=True, )从实现上可以补充两点对理解至关重要的机制:
- kwargs 通过 pickle 协议跨语言传递:Python 侧的 kwargs 会经
pickle.dumps(kwargs, protocol=5)序列化为字节串(见 plugins.py 的_serialize_kwargs,其中指出协议 5 是serde-picklecrate 支持的最高协议),Rust 侧再反序列化为自定义 struct。因此你的 kwargs 值必须可被 pickle 序列化; - kwargs 也参与 schema(字段类型)推导:引擎调用
plugin_field时会一并传入 kwargs 字节串(见下节与 plugin.rs 的minor == 1分支),所以 kwargs 可以影响输出列的 schema 计算。
仓库内完整示例 pyo3-polars/example/derive_expression/expression_lib/src/expressions.rs 中还包含一个更贴近真实场景的 kwargs 用例:change_time_zone使用output_type_func_with_kwargs,在输出类型函数convert_timezone中读取 kwargs 里的时区字符串并把Datetime的输出 dtype 改为带时区类型。
输出数据类型:让输出类型跟随输入类型变化
插件的输出数据类型不一定是固定的,往往取决于输入列的类型。为此,#[polars_expr()]宏支持output_type_func参数,指向一个把输入字段&[Field]映射为输出Field(列名 + 数据类型)的函数。polars_plan::dsl::FieldsMapper提供了常见映射的工具化封装。
下面的例子实现 haversine(球面距离)计算:输入为四列经纬度浮点数,输出希望保持与输入相同的浮点精度(Float32输入出Float32,Float64输入出Float64),因此输出类型无法写死,而需要由输入字段动态决定:
use polars_plan::dsl::FieldsMapper; fn haversine_output(input_fields: &[Field]) -> PolarsResult<Field> { FieldsMapper::new(input_fields).map_to_float_dtype() } #[polars_expr(output_type_func=haversine_output)] fn haversine(inputs: &[Series]) -> PolarsResult<Series> { let out = match inputs[0].dtype() { DataType::Float32 => { let start_lat = inputs[0].f32().unwrap(); let start_long = inputs[1].f32().unwrap(); let end_lat = inputs[2].f32().unwrap(); let end_long = inputs[3].f32().unwrap(); crate::distances::naive_haversine(start_lat, start_long, end_lat, end_long)? .into_series() } DataType::Float64 => { let start_lat = inputs[0].f64().unwrap(); let start_long = inputs[1].f64().unwrap(); let end_lat = inputs[2].f64().unwrap(); let end_long = inputs[3].f64().unwrap(); crate::distances::naive_haversine(start_lat, start_long, end_lat, end_long)? .into_series() } _ => polars_bail!(InvalidOperation: "only supported for float types"), }; Ok(out) }要点:
#[polars_expr(output_type_func=haversine_output)]会把输出类型的推导委托给haversine_output,引擎在进行 schema 推导、优化与生成物理计划时都会调用它;FieldsMapper::map_to_float_dtype()会把输出 dtype 映射为输入浮点 dtype;函数体内再按Float32/Float64两个分支分别执行核函数,其他类型直接polars_bail!报错——这种「schema 函数 + 分派实现」的组合是处理多态输入的标准套路;- 宏还支持
output_type_func_with_kwargs(输出类型同时依赖 kwargs)与固定output_type=DataType(输出类型恒定)三种模式,关键字定义见 pyo3-polars/pyo3-polars-derive/src/keywords.rs。
在多输入、需要统一输入类型的场景,Python 侧可以在注册时设置cast_to_supertype=True,让 Polars 先把各输入列 cast 到公共超类型再交给插件——仓库示例 dist.py 中注册四输入haversine时即使用该选项。
register_plugin_function 完整参数语义
register_plugin_function的行为参数直接决定 Polars 引擎如何处理你的函数,声明错误会带来错误结果甚至崩溃,其完整签名与语义如下(源码见 py-polars/src/polars/plugins.py,Rust 侧对应绑定签名见 py-polars/src/polars/_plr.pyi):
| 参数 | 语义 |
|---|---|
plugin_path | 插件包路径。接受动态库文件的直接路径,或包含动态库的目录路径(会扫描其中.so/.dll/.pyd)。路径默认相对于 Python 虚拟环境解析,可通过use_abs_path=True强制按绝对路径解析 |
function_name | 要注册的 Rust 函数名,必须与#[polars_expr]函数名完全一致 |
args | 传给函数的一列或多列表达式(IntoExpr或IntoExpr迭代器),对应 Rust 侧的inputs参数 |
kwargs | 非表达式参数(必须是可 JSON / pickle 序列化的普通值) |
is_elementwise | 声明函数仅对每个标量独立操作,可能触发快速路径并允许批量/并行执行 |
changes_length | 声明函数会改变表达式长度,例如unique、slice这类操作 |
returns_scalar | 当函数作为最终聚合运行且输出为单元长度时自动 explode,适用于sum、min、covariance等聚合语义 |
cast_to_supertype | 调用前先把输入表达式 cast 到公共超类型 |
input_wildcard_expansion | 在执行函数前展开通配符表达式(如pl.col("*")) |
pass_name_to_apply | 设为True时,在 group-by 中传给函数的 Series 会保证列名被正确设置(每组多一次堆分配) |
use_abs_path | 为True时把plugin_path解析为绝对路径(默认为相对虚拟环境的路径) |
底层原理:引擎如何动态加载并调用你的插件
理解插件的 C ABI 约定有助于排查问题(例如「符号未找到」与版本不匹配)。Polars 主引擎在 crates/polars-plan/src/plans/aexpr/function_expr/plugin.rs 中实现了全部加载逻辑:
- 缓存动态库:
LOADED(LazyLock<RwLock<PlIndexMap<String, Arc<PluginAndVersion>>>>)以库路径为键缓存已加载的Library。重复调用不会重复dlopen; - 定位并加载:Python 构建下,相对路径会基于
sys.prefix(虚拟环境根)拼接为绝对路径后再交给libloading::Library::new打开;加载失败会包装为ComputeError; - 版本握手:加载后立即查找符号
_polars_plugin_get_version,读出 32 位版本号并拆分为major与minor(高 16 位 / 低 16 位)。若major != 0,引擎会以「此 Polars 引擎不支持该插件版本」为由拒绝执行; - 字段/schema 推导:调用
plugin_field,通过libloading查找导出符号_polars_plugin_field_{fn_name},把输入字段序列化为 Arrow C 结构(ArrowSchema)传给插件,插件返回输出字段(可附带 kwargs 字节串); - 执行:
call_plugin通过符号_polars_plugin_{fn_name}调用插件函数。输入列经由polars_ffi的export_column封装为SeriesExport,连同 kwargs 字节串与默认CallerContext一起传入,返回值写回SeriesExport; - 错误与 panic 传播:若返回值为空,引擎通过导出符号
_polars_plugin_get_last_error_message读取线程局部错误字符串。插件侧发生 panic 时,派生宏会用std::panic::catch_unwind捕获并标记为"PANIC";引擎检测到后抛错并提示可设置POLARS_VERBOSE=1将 panic 信息输出到 stderr(见 plugin.rs 的check_panic)。
对应的符号生成规则在派生宏侧可以找到:#[polars_expr]宏把你的函数编译为#[no_mangle] pub unsafe extern "C"导出,符号名分别按_polars_plugin_{fn_name}与_polars_plugin_field_{fn_name}生成,参见 pyo3-polars/pyo3-polars-derive/src/lib.rs 的get_expression_function_name与get_field_function_name。这解释了为什么文档反复强调函数名必须拼写正确——一旦与宏生成的符号不一致,引擎将找不到入口。
源码级进阶:并行、日期与更多示例函数
仓库中的官方示例工程不只是 Pig Latin 的最小实现,它还覆盖了若干进阶模式,值得作为模板研读:
- 手动并行分块:
pig_latinnify_with_parallelism接收CallerContext作为参数(位于 Rust 函数第二参数位置,可与 kwargs 同时出现),在context.parallel()为真时把字符串列切分到多个线程分块处理再重组,展示了如何配合引擎执行上下文主动引入rayon并行。切分逻辑见 expressions.rs 中的split_offsets; - 多输入 elementwise:
hamming_distance使用arity::binary_elementwise_values对两个字符串列逐对计算汉明距离;jaccard_similarity则用arity::binary_elementwise对两个整数列表列计算 Jaccard 相似度,并调用polars_ensure!做输入类型前置校验(见 distances.rs); - 日期类型处理:
is_leap_year通过input.date()取DateChunked,借助as_date_iter()遍历可选日期并调用dt.leap_year(),收集为BooleanChunked; - 调用方约束:pyproject.toml 示例(pyo3-polars/example/derive_expression/expression_lib/pyproject.toml)声明以
maturin>=1.0,<2.0作为构建后端;工程根目录的 Makefile 与 run.py 提供了编译与端到端运行脚本,可直接执行验证。
结语:何时选择表达式插件
当遇到以下情况时,表达式插件是最优解:需要把一段 Rust 算法(尤其是逐元素、多列协同或分组内部逻辑)变成可与原生表达式混用的Expr;对单次调用性能敏感且希望完全避开 Python 与 GIL;或希望复用既有的 Rust 计算内核。而如果你的需求只是小型脚本内的临时逻辑、不需要极致性能与并行,先尝试内置表达式组合往往更快落地。插件需要为调用方(Python 包名)、Rust 函数名、ABI 版本三者保持一致负责,本文给出的源码级调用链与官方示例工程将帮助你建立这一完整心智模型,从而写出可维护、可并行的 Polars 自定义表达式。
【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考