十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Polars 插件(Plugins)扩展机制全指南:表达式插件与 IO 插件的原理与实战

Polars 插件(Plugins)扩展机制全指南:表达式插件与 IO 插件的原理与实战 Polars 插件Plugins扩展机制全指南表达式插件与 IO 插件的原理与实战【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polarsPolars 内置了大量表达式与数据源能力但真实业务中总有引擎覆盖不到的自定义算法与私有文件格式。为此Polars 在官方用户手册中专门开辟了plugins一节系统性地给出「表达式插件Expression plugins」与「IO 插件IO plugins」两条官方扩展路线并整理了社区插件与配套学习资料。本文以 plugins/index.md 为主线结合仓库内pyo3-polars派生宏、Python 侧注册入口与可运行示例的源码实现完整讲解如何在 Polars 中编写、注册并分发自己的插件读完即可照例复刻出一个编译期链接、性能接近原生表达式的自定义函数或一个接入懒执行/流式引擎的私有数据源。插件体系的整体蓝图一份「入口索引」读懂两条扩展路线Polars 官方手册的 plugins 一节本质是一张扩展能力地图它把插件分为互补的两类表达式插件Expression plugins把一段编译好的 Rust 函数注册成 Polars 表达式适合扩展计算内核IO 插件IO plugins把自定义的读取源注册进查询引擎适合接入 Polars 尚不支持的输入格式。两者详细指南分别见 Expression plugins 与 IO plugins它们共同构成plugins/index.md这一入口页的全部主线内容。理解这个「入口页 → 两篇子指南」的层次也就理解了插件系统的设计分界凡是「一行输入产生一行输出或一组聚合的计算」走表达式插件凡是「从某种介质把数据灌进引擎」的读取走 IO 插件。社区插件生态与配套学习资料原索引页内容入口页还维护了一份精选的、非穷举curated, non-exhaustive的社区插件清单可作为扩展方向的参考以下插件均为仓库外项目具体可用性以各项目为准综合类polars-xdt补充主库暂未纳入的日期时间类功能、polars-hash稳定的非加密与加密哈希函数数据科学类polars-distance成对距离函数、polars-ds面向常规数值/字符串分析流程的简化扩展生物信息类polars-bio构建在 Polars、Arrow、DataFusion 之上的基因组学 Python 库地理空间类polars-st在 DataFrame/Series/Expr 上提供类似 Shapely、Geopandas 的空间运算、polars-reverse-geocode离线逆地理编码、polars-h3接入 H3 离散全球网格系统把点与几何直接索引为六边形。入口页同时给出了三类补充材料官方作者的插件主题演讲Ritchie Vink Keynote on Polars Plugins、polars-plugins-tutorial极简教程以及cookiecutter-polars-plugin项目脚手架模板。若想在本仓库内直接看到「完整可编译」的对应样例推荐优先研读两个官方示例目录pyo3-polars/example/derive_expression表达式插件完整工程与 pyo3-polars/example/io_pluginRust 实现的 IO 源后文所有代码讲解都能在这两处找到一一对应的落地实现。表达式插件接近原生性能的自定义函数表达式插件是 Polars 实现用户自定义函数UDF的首选方式。它把一个编译好的 Rust 函数注册为 Polars 表达式查询引擎在运行时动态链接该函数执行速度几乎与原生表达式相当整个过程完全不经过 Python、不产生 GIL 竞争。作为一等公民自定义表达式自动继承原生表达式的三类收益查询优化Optimization可被谓词下推、投影下推、常量折叠等优化规则感知与重排并行执行Parallelism由引擎调度到多线程Rust 原生性能Rust native performance无解释器开销。文档以「Pig Latin 转换」作为第一个教学案例——把每个单词的首字母挪到词尾并追加ay如pig→igpay。该功能当然能用纯表达式拼接col(name).str.slice(1) col(name).str.slice(0, 1) ay实现但专函数化既更高效也是学习插件机制的最小切口。工程搭建Cargo.toml与编译目标先创建一个新的 Rust library其Cargo.toml关键点在于crate-type [cdylib]最终产物是供动态链接的共享库并依赖polars、pyo3启用extension-module与abi3-py310以及带derivefeature 的pyo3-polars提供#[polars_expr]过程宏示例如下[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] }仓库内的真实样例 pyo3-polars/example/derive_expression/expression_lib/Cargo.toml 结构与此一致并且其lib.rs还会安装一个全局内存分配器use pyo3_polars::PolarsAllocator; mod distances; mod expressions; #[global_allocator] static ALLOC: PolarsAllocator PolarsAllocator::new();如 pyo3-polars/example/derive_expression/expression_lib/src/lib.rs 所示安装PolarsAllocator是为保证插件内 Rust 分配的缓冲区与 Polars 的分配机制保持一致避免跨库free引发内存错误——编写真实插件时这一步值得保留。Rust 侧#[polars_expr]宏与 Series 处理定义插件函数的核心约束有两条必须给函数标注#[polars_expr(output_typeDataType)]函数签名必须把inputs: [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_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()) }这里特意选用apply_into_string_amortized而非apply_values是为了复用同一个输出缓冲避免每行都新建 String。若插件接受多个输入、按元素操作并产出String文档建议关注polars::prelude::arity下提供的binary_elementwise_into_string_amortized等工具函数它们按同样的零额外分配思路处理二元输入。#[polars_expr]宏在编译期替你生成了跨 FFI 的接线代码。从派生宏实现 pyo3-polars/pyo3-polars-derive/src/lib.rs 可以看清其底层契约宏会把你的函数调用包装成fn(inputs, [kwargs] [, context]) - PolarsResultSeries成功时通过polars_ffi::version_0::export_series(out)把Series导出给 Polars失败时则调用_update_last_error记录错误并返回空值——这也就是文档所说的「引擎在运行时动态链接」的真实载体。Python 侧包结构、注册与调用Rust 侧到此即全部完成。Python 侧需要一个与Cargo.toml中[lib] name完全同名的文件夹本例为expression_lib与 Rust 的src平级放置里面至少有一个__init__.py。最终目录结构如下├── expression_lib/ # name must match lib.name in Cargo.toml │ └── __init__.py │ ├── src/ │ ├── lib.rs │ └── expressions.rs │ ├── Cargo.toml └── pyproject.toml注册函数时function_name必须与 Rust 侧函数名完全一致否则主 Polars 包无法解析同时可通过关键字参数向引擎声明函数行为。下面声明is_elementwiseTrue表示函数是逐元素操作——只有这样才能允许 Polars 按批次执行像sort、slice这类会改变数据位置或长度的操作则绝不能如此标记。# 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, )注意这里PLUGIN_PATH Path(__file__).parent指向插件包目录本身Python 侧注册入口在 py-polars/src/polars/plugins.py 中定义。仓库里真实示例的注册方式如 expression_lib/language.py 所示——把「指向.so的位置」与「函数名」显式分离def pig_latinnify(expr: IntoExprColumn, capitalize: bool False) - pl.Expr: return register_plugin_function( plugin_pathLIB, args[expr], function_namepig_latinnify, is_elementwiseTrue, kwargs{capitalize: capitalize}, )register_plugin_function的行为声明参数Python 注册函数的行为声明直接决定 Polars 引擎如何处理你的函数官方实现py-polars/src/polars/plugins.py支持以下参数它们的含义即为引擎优化与调度所依赖的「事实」参数作用plugin_path插件包路径可指向动态库文件本身或其所在目录默认相对于虚拟环境解析设置use_abs_pathTrue可改为绝对路径function_name要注册的 Rust 函数名必须与 Rust 侧严格一致args传给函数的表达式参数对应 Rust 侧inputs可传单个表达式或表达式迭代器kwargs非表达式参数如浮点/整型/字符串/布尔需能被序列化经 FFI 到 Rust 侧由 serde 反序列化is_elementwise声明函数只对标量逐元素操作可能触发快路径changes_length声明函数会改变表达式长度如unique、slicereturns_scalar若函数作为最终聚合运行自动对单元素结果做explode如sum、min等聚合语义cast_to_supertype执行前把多个输入表达式转型到公共超类型input_wildcard_expansion在函数执行前展开*通配表达式pass_name_to_apply在 group-by 场景下确保传给函数的 Series 名称被正确设置每个分组多一次堆分配源码 docstring 同时给出严厉警告注册参数若与真实行为不符会导致错误结果因此这些开关必须与 Rust 实现语义严格对应。编译与首次运行在当前环境安装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))如果你希望把它收进一个更像原生语法的命名空间也可以注册自定义 namespaceregister_expr_namespace让用户写出out df.with_columns( pig_latinpl.col(convert).language.pig_latinnify(), )即注册一个Expr.language命名空间把多个相关插件函数统一收纳。进阶一通过 kwargs 传递自定义参数若插件需要接收普通标量参数只需在 Rust 侧定义一个struct并派生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()) }关键点有二kwargs必须位于函数签名第二个位置struct 的字段类型会决定 Python 端必须传什么。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, )仓库示例 expression_lib/language.py 中的append_args与 Rust 侧append_kwargs一一对应是一份可编译对照的现成参考。派生宏内部对 kwargs 的处理也可见于 pyo3-polars-derive/src/lib.rs它通过_parse_kwargs完成反序列化失败时返回带提示语的InvalidOperation错误。进阶二根据输入类型动态决定输出类型输出类型并不总是固定的常常取决于输入表达式类型。此时可给#[polars_expr()]提供output_type_func参数指向一个把输入字段[Field]映射为输出Field含列名与数据类型的函数。典型做法是借助工具结构FieldsMapperuse 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) }上例让输出类型跟随输入在Float32/Float64间变化遇到其他类型则以polars_bail!报错。其底层实现haversine 距离计算等在仓库示例 pyo3-polars/example/derive_expression/expression_lib/src/distances.rs 与expressions.rs中均有完整实现读者可对照output_type_func与函数体的「类型分派」写法理解这套双轨机制。IO 插件为懒执行引擎接入自定义数据源除表达式插件外Polars 还支持IO 插件——把自定义文件格式注册为查询引擎的数据源source。IO 源可以通过 Arrow FFI 以零拷贝方式把大块数据交给引擎因此官方选择现阶段用 Python 作为 IO 插件的接口层引擎认为 GIL 仅在短暂的数据交接点被持有不会成为瓶颈。例如一个 IO 源可以在 Rust 中解析自己的 DataFrame仅在 rendez-vous 时刻零拷贝移交数据、短时占用 GIL。适用场景当你有 Polars 原生不支持的源文件却想继续享受投影下推、谓词下推、提前终止early stopping与流式引擎streaming engine等优化时就该考虑 IO 插件。教学示例一个刻意写得很慢的自定义 CSV 源文档强调以下示例仅为教学目的性能刻意做差。首先准备导入其中关键的是register_io_source每次实例化时注册一个新的生成器# Use python for csv parsing. import csv import polars as pl # Used to register a new generator on every instantiation. from polars.io.plugins import register_io_source from typing import Iterator import io第一步解析 schemaPolars 中每个scan函数都必须能提供其读取数据的 schema。这个简化版 CSV 解析器把所有列一律读成pl.String差异只在字段名与字段个数def parse_schema(csv_str: str) - pl.Schema: first_line csv_str.split(\n)[0] return pl.Schema({k: pl.String for k in first_line.split(,)})对小型 CSVa,b,c\n1,2,3运行将得到Schema([(a, String), (b, String), (c, String)])。第二步写数据源外层包装 内层生成器IO 源采用「外函数 内函数」双层结构。外层my_scan_csv是面向用户的入口接收文件名等自定义参数CSV 场景下如delimiter、quote_char内层才是数据源实现其签名是预定义、必须接受的共四个参数with_columns被投影的列。若应用了投影下推读取端应尽量只产出这些列predicatePolars 表达式。读取端应据此过滤行n_rows只需物化前 n 行读取端读完即可停止batch_size期望的批大小提示读取端生成器应尽量按此产块。def my_scan_csv(csv_str: str) - pl.LazyFrame: schema parse_schema(csv_str) def source_generator( with_columns: list[str] | None, predicate: pl.Expr | None, n_rows: int | None, batch_size: int | None, ) - Iterator[pl.DataFrame]: Generator function that creates the source. This function will be registered as IO source. if batch_size is None: batch_size 100 # Initialize the reader. reader csv.reader(io.StringIO(csv_str), delimiter,) # Skip the header. _ next(reader) # Ensure we dont read more rows than requested from the engine while n_rows is None or n_rows 0: if n_rows is not None: batch_size min(batch_size, n_rows) rows [] for _ in range(batch_size): try: row next(reader) except StopIteration: n_rows 0 break rows.append(row) df pl.from_records(rows, schemaschema, orientrow) n_rows - df.height # If we would make a performant reader, we would not read these # columns at all. if with_columns is not None: df df.select(with_columns) # If the source supports predicate pushdown, the expression can be parsed # to skip rows/groups. if predicate is not None: df df.filter(predicate) yield df return register_io_source(io_sourcesource_generator, schemaschema)最外层调用register_io_source返回一个真正的pl.LazyFrame——于是它自动被并入懒执行计划享受下推与流式执行。第三步跑一个很慢的实测用两组字符串模拟「文件内容」验证完整 collect 与「头部截断 谓词下推」两种路径csv_str1 a,b,c,d 1,2,3,4 9,10,11,2 1,2,3,4 1,122,3,4 print(my_scan_csv(csv_str1).collect()) csv_str2 a,b 1,2 9,10 1,2 1,122 print(my_scan_csv(csv_str2).head(2).collect())运行输出分别如下第一个完整打印 4 行 4 列第二个受.head(2)驱动内部以n_rows截断只产出了前 2 行这正是文档所述「reader can stop when n_rows are read」的体现shape: (4, 4) ┌─────┬─────┬─────┬─────┐ │ a ┆ b ┆ c ┆ d │ │ --- ┆ --- ┆ --- ┆ --- │ │ str ┆ str ┆ str ┆ str │ ╞═════╪═════╪═════╪═════╡ │ 1 ┆ 2 ┆ 3 ┆ 4 │ │ 9 ┆ 10 ┆ 11 ┆ 2 │ │ 1 ┆ 2 ┆ 3 ┆ 4 │ │ 1 ┆ 122 ┆ 3 ┆ 4 │ └─────┴─────┴─────┴─────┘ shape: (2, 2) ┌─────┬─────┐ │ a ┆ b │ │ --- ┆ --- │ │ str ┆ str │ ╞═════╪═════╡ │ 1 ┆ 2 │ │ 9 ┌10 ┐ ... └─────┴─────┘注第二个 DataFrame 仅两行两列示意投影与行数限制均已生效。源码级细节register_io_source的完整签名与引擎约定Python 侧实现位于 py-polars/src/polars/io/plugins.py从它的签名与 docstring 可以提炼出官方对 IO 插件作者的约定接口本身被标记为unstable()即 API 可能在不视为破坏性变更的前提下随时调整生产接入需留意版本锁定io_source一个签名为(with_columns, predicate, n_rows, batch_size) - Iterator[DataFrame]的生成器函数schema接受一个SchemaDict或返回 schema 的可调用对象描述投影下推之前的完整输出 schemavalidate_schema是否让引擎校验每个批次与给定 schema 匹配——不匹配属于实现错误校验能避免难以排查的隐性 bugis_pure若源是纯的同参数同结果优化器可对计划中重复出现的同一 IO 源做去重de-duplicationexplain_name/explain_detail自定义explain查询计划中该 scan 节点的标签与细节说明渲染形如PYTHON[name] SCAN ...。引擎传入的predicate会被序列化为字节并尝试用pl.Expr.deserialize还原若还原失败插件会收到None由 Polars 引擎自行在解析后过滤见 io/plugins.py 中对parsed_predicate_success的处理。仓库佐证一个真正在 Rust 中实现的 IO 源若想看到「底层用 Rust 实现 IO 源」的完整形态仓库提供了官方示例 pyo3-polars/example/io_plugin/io_plugin/src/lib.rs。它定义了一个RandomSource一个#[pyclass]其接口恰好映射引擎的四项约定schema()返回PySchema让引擎拿到完整输出 schematry_set_predicate(predicate: PyExpr)与set_with_columns(columns: VecString)分别承接谓词下推与投影下推next()每次按size_hint/n_rows产出一个PyDataFrame内部用注释明确区分「投影下推/切片下推应在采样前做」「谓词下推可以在更低层完成、也可事后过滤」——与 Python 教学版的两阶段处理一一对应。对应 Python 侧封装在 pyo3-polars/example/io_plugin/io_plugin/io_plugin/init.pyscan_random先创建采样器再在生成器中把with_columns、predicate透传给 Rust 对象若 Rust 端因无法反序列化谓词而返回失败标记predicate_set False则退回在 Python 端out.filter(predicate)。这套「先尝试下推、失败则本地过滤」的兜底逻辑是编写高性能 IO 源时值得复用的模式。如何选型表达式插件还是 IO 插件综合两条路线可归纳出清晰的选型建议诉求选型关键收益在列/元素粒度扩展计算能力字符串、数学、哈希、地理、生物信息算法等表达式插件无 Python/GIL 开销共享优化/并行/原生性能语义开关由注册参数显式声明接入新文件格式/数据源并享受懒执行优化IO 插件零拷贝 Arrow FFI、投影/谓词下推、提前终止、流式引擎支持现阶段以 Python 为接口层只是想快速体验直接阅读并运行 derive_expression 示例 与 io_plugin 示例两者均带Makefile与run.py可编译、可运行的最小闭环两条路线的权威说明分别见 expr_plugins.md 与 io_plugins.md而 plugins/index.md 作为总入口负责把读者导流到对应方向并呈现社区生态。需要特别提醒的事实边界是入口页所列举的社区插件均为仓库外独立项目具体能力与维护状态以其各自仓库为准本仓库内可直接验证的仅是官方示例工程与 Python/Rust 侧的注册、宏展开与 FFI 实现。对想动手写第一个插件的开发者最经济的路径是先跑通derive_expression示例的make/run.py再对照本文把 Pig Latin 与 kwargs 两个案例亲手改一遍当需要处理私有格式时再基于io_plugin示例把RandomSource替换为真实解析器即可。【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表