Polars Expression Plugins 完全指南:用 Rust 编写原生速度的 DataFrame 表达式插件

发布时间:2026/9/9 19:56:33
Polars Expression Plugins 完全指南:用 Rust 编写原生速度的 DataFrame 表达式插件 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_typeDataType)]属性且第一个参数必须是inputs: [Series]返回PolarsResultSeries// 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_typeString)] fn pig_latinnify(inputs: [Series]) - PolarsResultSeries { 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_pathPLUGIN_PATH, function_namepig_latinnify, argsexpr, is_elementwiseTrue, )其中几个参数必须理解到位function_name必须与 Rust 侧的#[polars_expr]函数名完全一致此处均为pig_latinnify否则主 Polars 包无法解析到对应符号plugin_path指向插件包所在目录运行时会被解析为其中.so/.dll/.pyd动态库解析逻辑见 py-polars/src/polars/plugins.py 中的_resolve_plugin_path与_is_dynamic_libis_elementwiseTrue告知 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_latinpig_latinnify(convert))进阶把插件挂载为自定义命名空间除了「函数式调用」外还可以通过 Polars 的命名空间注册 API 创建自定义命名空间例如Expr.language让用户以链式调用的风格编写out df.with_columns( pig_latinpl.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_typeString)] fn append_kwargs(input: [Series], kwargs: MyKwargs) - PolarsResultSeries { 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_pathPLUGIN_PATH, function_nameappend_kwargs, argsexpr, kwargs{ float_arg: float_arg, integer_arg: integer_arg, string_arg: string_arg, boolean_arg: boolean_arg, }, is_elementwiseTrue, )从实现上可以补充两点对理解至关重要的机制kwargs 通过 pickle 协议跨语言传递Python 侧的 kwargs 会经pickle.dumps(kwargs, protocol5)序列化为字节串见 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输入出Float32Float64输入出Float64因此输出类型无法写死而需要由输入字段动态决定use polars_plan::dsl::FieldsMapper; fn haversine_output(input_fields: [Field]) - PolarsResultField { FieldsMapper::new(input_fields).map_to_float_dtype() } #[polars_expr(output_type_funchaversine_output)] fn haversine(inputs: [Series]) - PolarsResultSeries { 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_funchaversine_output)]会把输出类型的推导委托给haversine_output引擎在进行 schema 推导、优化与生成物理计划时都会调用它FieldsMapper::map_to_float_dtype()会把输出 dtype 映射为输入浮点 dtype函数体内再按Float32/Float64两个分支分别执行核函数其他类型直接polars_bail!报错——这种「schema 函数 分派实现」的组合是处理多态输入的标准套路宏还支持output_type_func_with_kwargs输出类型同时依赖 kwargs与固定output_typeDataType输出类型恒定三种模式关键字定义见 pyo3-polars/pyo3-polars-derive/src/keywords.rs。在多输入、需要统一输入类型的场景Python 侧可以在注册时设置cast_to_supertypeTrue让 Polars 先把各输入列 cast 到公共超类型再交给插件——仓库示例 dist.py 中注册四输入haversine时即使用该选项。register_plugin_function 完整参数语义register_plugin_function的行为参数直接决定 Polars 引擎如何处理你的函数声明错误会带来错误结果甚至崩溃其完整签名与语义如下源码见 py-polars/src/polars/plugins.pyRust 侧对应绑定签名见 py-polars/src/polars/_plr.pyi参数语义plugin_path插件包路径。接受动态库文件的直接路径或包含动态库的目录路径会扫描其中.so/.dll/.pyd。路径默认相对于 Python 虚拟环境解析可通过use_abs_pathTrue强制按绝对路径解析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 中实现了全部加载逻辑缓存动态库LOADEDLazyLockRwLockPlIndexMapString, ArcPluginAndVersion以库路径为键缓存已加载的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_VERBOSE1将 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多输入 elementwisehamming_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声明以maturin1.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),仅供参考