使用 Rust 表达式插件扩展 Polars
摘要
本文解释了 fenic 为何以及如何使用 Rust 表达式插件来扩展 Polars,以在引擎原生执行文本操作(分块、提示模板化、模糊匹配等),从而避免 Python UDF 的性能和组合问题。
暂无内容
查看缓存全文
缓存时间: 2026/07/24 20:03
# 我们为何以及如何用 Rust 表达式插件扩展 Polars,服务于 fenic | fenic 博客
来源:https://fenic.ai/blog/extending-polars-with-rust-expression-plugins
fenic 是一个语义级 DataFrame 库。它提供 PySpark 风格的 API,用于在混乱的非结构化数据上构建 AI 和 LLM 管道。其本地引擎是 Polars。本文讲的是我们在构建它时遇到的一个具体问题,以及 Polars 中解决该问题的部分:表达式插件系统。我们最终编写了九个 Rust 插件来扩展 Polars 的表达式引擎。以下内容既包括我们选择这条路径的原因,也包括如何使用真实代码构建它们。
**tl;dr.** AI 管道在文本上需要的操作(分块、提示模板、jq、模糊匹配、Markdown 和转录解析、更丰富的类型转换)并不在 Polars 中。把它们作为 Python UDF 来实现很慢,而且会破坏组合性。通过 `pyo3-polars`(https://github.com/pola-rs/pyo3-polars)将它们编写为 Rust 的 Polars 表达式插件,就能把它们变成原生的表达式。它们在引擎内对 Arrow 运行,保持声明的类型,并且可以在单个表达式树中与内置操作组合。如果你在 UDF 和插件之间权衡,这就是选择插件的论据。
---
## 为什么我们要构建这个
### fenic 是什么
fenic 是一个用于构建 AI 和 LLM 管道的 DataFrame 库,API 模仿 PySpark。如果你来自 Polars 背景,需要知道的重要事情是:Polars 是 fenic 的核心执行引擎。你编写的每一个 DataFrame 操作都会变成 Polars 表达式,或者基于 `pl.DataFrame` 的计划。语义操作、LLM 的 map/extract/classify、嵌入、相似性连接,所有这些都建立在同样的机制之上。
因此,实际上 fenic 能做的事情受限于我们能在 Polars 中表达的内容。这是一个刻意的选择。Polars 为我们提供了一个快速、列式、Arrow 原生的引擎,带有真正的表达式语言和惰性优化器。我们对重新构建其中的任何部分都没有兴趣。我们只是想在其之上添加功能。
### 我们想要做什么
fenic 目标的工作负载是处理非结构化文本的管道。文档、聊天记录、转录文本、抓取的 JSON、Markdown。具体来说,这意味着 Polars 没有原生等价物的行级操作:
- **分块**:将文档按*token*计数分成交叠窗口,用于嵌入或检索。
- **渲染提示模板**:每行一个。真正的 Jinja,以列结构体作为变量,在 LLM 调用之前使用。
- **使用 jq 查询 JSON** 列,以及**解析 Markdown** 为结构化 AST。
- **模糊匹配**字符串(六种编辑距离指标),用于去重和连接。
- **解析转录文本**(SRT/WebVTT)为带时间戳的、类型化的提示记录。
- **将值转换**为 fenic 更丰富的逻辑类型(嵌入、Markdown、类型化结构体),这些类型是 Polars 的物理 dtype 无法直接建模的。
以上每一种操作都必须对整个列运行,产生一个类型化的结果,并嵌入到更大的 DataFrame 管道的中间,而不是在旁边独立运行。
### 显而易见的方法在哪些方面不足
向 Polars 添加自定义操作的直观方式是 Python UDF。`map_elements` 用于逐行工作,`map_batches` 用于整个 Series 的工作。对于真正不透明、受 IO 限制的工作(比如 LLM API 调用),这仍然是正确的工具,fenic 也正是为此使用 `map_batches`。但对于上述文本操作,UDF 在三个方面代价过高。
**速度。** `map_elements` 在 GIL 下每行运行 Python,每次值都需要 Python 对象往返。对于对数百万行进行 tokenize 或模糊匹配,这就是瓶颈——不是工作本身。
**组合性。** 这是真正让人头疼的一点。Polars 之所以快,是因为它整个规划和执行表达式树。一旦某一步是一个不透明的 Python 回调,引擎就无法看透它。它变成了优化障碍,强制物化,破坏了单次通过的管道。将三个 UDF 链起来,就会产生三次进出引擎的往返。
**类型。** UDF 的输出类型是你松散断言并希望成立的。我们需要能产生真实、声明的 dtype 的操作,比如 `List`、`Struct`、固定大小的嵌入数组,这样计划的其余部分可以针对它们进行类型检查。
其他替代方案更糟糕。Fork Polars 来添加原生内核意味着永远维护一个 fork。在 DataFrame 外部进行文本预处理,然后加载,这扔掉了惰性和组合性——而这正是使用 Polars 的全部原因。我们希望这些操作成为表达式引擎中的一等公民,而不是它的邻居。
### 为什么选择插件
Polars 有一个专门用于此目的的内置答案:表达式插件。你用 Rust 编写一个内核,注册它,它就变成了一个普通的 `pl.Expr`。对引擎来说,它与内置函数无法区分。这满足了 UDF 路径错过的所有条件。
它在引擎内、Rust 中、直接操作 Arrow 缓冲区运行。向量化、可并行、可流式处理,没有 Python,也没有逐行对象往返。它返回一个 `pl.Expr`,因此可以组合。插件可以相互链接,也可以与原生操作一起组合成引擎整体规划的单表达式树。它声明输出 dtype,因此计划在自定义步骤期间保持类型化。而且 Rust 为我们提供了成熟的热循环生态系统(`jaq`、`minijinja`、`rapidfuzz`、`tiktoken`、Markdown 解析器),无需重新实现任何东西。
我们定下的规则不是“把所有东西都用 Rust 重写”。原生 Polars 仍然是快速路径。插件只填补 Polars 无法表达的空缺:一个分词器、一个 jq 引擎、一个捕获组索引本身就是一个列的 regex。即使是模糊匹配器,也只将六个原始内核保留在 Rust 中,并从普通的 Polars 表达式组合出更高级的比率。`pyo3-polars` 生成 FFI、Arrow 编组以及关键字参数桥接,因此编写一个插件的成本很低。我们写了九个。
### 它给我们带来了什么
在讨论机制之前,值得先说明最终状态。fenic 的每一个文本操作现在都是原生的 Polars 表达式。它们在引擎内对共享的 Arrow 内存运行,保持它们的类型,而最重要的是:它们与原生 Polars 操作可以在单个表达式中组合,零 Python 往返。解析一个 Markdown 文档,用 jq 过滤其 AST,索引结果,并将其转换为类型化结构体,就是一个表达式。Polars 在一个单次通过中规划和执行它,自定义 Rust 和内置的列表操作并肩工作。一旦上面所有部件都准备好了,我们会回到这个具体的例子。
本文的其余部分是关于它是如何构建的,从 Python 注册到 Arrow 内存,只使用真实代码。
---
## 它是如何构建的:实现演练
### 思维模型:围绕契约的两个薄层
一个 Polars 表达式插件由两小段代码组成,它们之间有一个定义明确的契约:
1. **Python 端。** 注册一个函数,使其看起来像原生 Polars:`expr.my_namespace.my_op(...)`。
2. **Rust 端。** 一个接受 `&[Series]`、返回 `PolarsResult` 并声明其输出 dtype 的函数。
两者之间的所有内容都由 `pyo3-polars` 为你生成:通过 Arrow 将 `Series` 跨越 FFI 边界传递、编组关键字参数、连接符号查找。你永远不需要手动接触 FFI。
以下是 fenic 的 `json.jq` 操作的所有 Python 表面积:
```python
# src/fenic/_backends/local/polars_plugins/json.py
from pathlib import Path
import polars as pl
from polars.plugins import register_plugin_function
PLUGIN_PATH = Path(__file__).parents[3]
@pl.api.register_expr_namespace("json")
class Json:
"""Namespace for JSON-related operations on Polars expressions."""
def __init__(self, expr: pl.Expr) -> None:
self.expr = expr
def jq(self, query: str) -> pl.Expr:
return register_plugin_function(
plugin_path=PLUGIN_PATH,
function_name="jq_expr",
args=self.expr,
kwargs={"query": query},
is_elementwise=True,
)
```
两个 Polars API 做了主要工作:
- `@pl.api.register_expr_namespace("json")` 把 `.json` 访问器挂接到进程中*每一个* `pl.Expr` 上。导入后,在任何 Polars 表达式合法的地方,`pl.col("payload").json.jq(".name")` 都是一个合法的表达式。
- `register_plugin_function(...)` 返回一个普通的 `pl.Expr`,当引擎求值时,它会 dlopen 位于 `plugin_path` 的编译库,查找名为 `function_name` 的符号,传入输入 `Series`,并读取结果。
这就是全部思路。`.json` 在 Polars 中没有任何特殊处理。它是一个用户注册的命名空间,fenic 注册了九个(`json`、`jinja`、`markdown`、`chunking`、`tokenization`、`fuzz`、`regexp`、`dtypes`、`transcript`)。插件看起来完全像内置函数,因为对于表达式引擎来说,没有有意义的区别。
(插件新手?请参考 Polars 插件文档(https://docs.pola.rs/user-guide/plugins/expr_plugins/)和 `pyo3-polars`(https://github.com/pola-rs/pyo3-polars)仓库作为权威参考。)
### 从头到尾演练一个插件
上面的 Python `jq` 方法指向一个名为 `jq_expr` 的 Rust 符号。以下是它在边界另一侧的完整代码:
```rust
// rust/src/json/mod.rs
use polars::prelude::*;
use polars_arrow::array::ValueSize;
use pyo3_polars::derive::polars_expr;
use serde::Deserialize;
use serde_json::Value;
#[derive(Deserialize)]
struct JqKwargs {
query: String,
}
fn jq_output(_: &[Field]) -> PolarsResult<Field> {
Ok(Field::new(
"jq".into(),
DataType::List(Box::new(DataType::String)),
))
}
#[polars_expr(output_type_func=jq_output)]
fn jq_expr(inputs: &[Series], kwargs: JqKwargs) -> PolarsResult<Series> {
// Compile the jq filter ONCE, reuse it for every row.
let filter = jq::build_jq_query(&kwargs.query)
.map_err(|e| PolarsError::ComputeError(e.to_string().into()))?;
let jq_inputs = RcIter::new(core::iter::empty());
let ca = inputs[0].str()?;
let mut builder = ListStringChunkedBuilder::new("jq".into(), ca.len(), ca.get_values_size() * 5);
for opt_str in ca.into_iter() {
if let Some(s) = opt_str {
match serde_json::from_str::<Value>(s) {
Ok(val) => {
let v: Val = val.into();
let results = filter
.run((Ctx::new([], &jq_inputs), v))
.collect::<Vec<_>>();
match results {
// an empty jq result set -> null
Ok(values) if values.is_empty() => builder.append_null(),
Ok(values) => {
let strs: Vec<String> = values.iter().map(|v| v.to_string()).collect();
builder.append_values_iter(strs.iter().map(|s| s.as_str()));
}
// a failed query aborts the whole batch — it does NOT null the row
Err(e) => return Err(PolarsError::ComputeError(
format!("jq query execution failed: {e}. Query: '{}'", kwargs.query).into(),
)),
}
}
// malformed JSON is unreachable: the column is typed JsonType, so
// upstream validation guarantees every non-null row is valid JSON.
Err(e) => unreachable!("Invalid JSON: {s} ({e})"),
}
} else {
builder.append_null();
}
}
Ok(builder.finish().into_series())
}
```
理解契约所需的一切都在这个代码片段中:
- **签名。** `#[polars_expr(...)]` 属性宏(来自 `pyo3_polars::derive`)将一个普通的 Rust 函数转换为导出的、C-ABI 的符号,Polars 可以 dlopen 并调用。你的函数只需要接收 `inputs: &[Series]` 并返回 `PolarsResult<Series>`。没有 `unsafe`,没有手动指针操作。宏生成 Arrow 编组。
- **类型化输入。** `inputs[0].str()?` 将第一个 `Series` 向下转换为 `StringChunked`。如果该列不是字符串,你会得到一个干净的 `PolarsError`,而不是段错误。
- **一次工作,而不是每行。** jq 过滤器只编译一次(`build_jq_query`),并在所有行之间重用。这是一个反复出现的模式:将编译、分配大小和设置提升出行循环之外。
- **使用构建器构建输出。** `ListStringChunkedBuilder`,预先用容量估计进行了大小调整,生成输出列。空值语义是有意为之。对于空输入行或空 jq 结果集,调用 `append_null()`。失败的 jq 查询会以 `ComputeError` 中止整个批次。而格式错误的 JSON 会触达 `unreachable!()`,因为输入列是 `JsonType` 类型,上游验证保证每个非空行都是有效 JSON。(这种保证是建立在类型化引擎上的一个不错副作用。内核可以假设其输入。)
- **返回一个 `Series`。** `builder.finish().into_series()`。`jaq` crate(纯 Rust jq)执行实际的过滤。Polars 提供列基础设施。插件是两者之间的接缝。
### 输出类型及其重要性
再看一下属性:`#[polars_expr(output_type_func=jq_output)]`。那个 `jq_output` 函数返回 `DataType::List(String)`,从不接触数据。这是关于插件需要正确理解的最重要的事情,也是最容易搞错的事情。
Polars 在执行表达式之前解析其 schema。插件是不透明的 Rust。引擎无法通过查看你的循环来推断输出。所以你必须事先声明输出 dtype,并且该声明必须在没有数据的情况下运行。
fenic 使用了宏支持的三种声明风格,根据输出类型依赖什么来选择。
**1. 静态。类型永远不会改变。** 所有六个模糊匹配内核总是返回一个 `Float64`:
```rust
// rust/src/fuzz/mod.rs
#[polars_expr(output_type=Float64)]
fn normalized_indel_similarity(inputs: &[Series]) -> PolarsResult<Series> { /* ... */ }
```
`tokenization.count_tokens` 同样是一个固定的 `#[polars_expr(output_type=UInt32)]`。
**2. 来自输入字段。类型是输入的函数。** `jq` 和文本分块总是产生 `List`,通过一个小函数声明,该函数接收输入 `Field`s(这里忽略它们)。
**3. 来自参数。类型取决于关键字参数。** 这是有趣的地方。fenic 的 `dtypes.cast` 内核实现了到 fenic 自己逻辑类型系统的转换,因此它的输出 dtype 实际上是一个参数,作为序列化的 JSON 类型描述符传递。输出类型函数在计划时反序列化该 JSON 以计算 Polars dtype:
```rust
// rust/src/dtypes/mod.rs
#[derive(Deserialize, Debug)]
pub struct CastKwargs {
source_dtype: String, // JSON-encoded fenic type
dest_dtype: String, // JSON-encoded fenic type
}
// Runs during schema inference only — never sees the data.
fn fenic_dtype_str_to_polars_dtype(
_input_fields: &[Field],
kwargs: CastKwargs,
) -> PolarsResult<Field> {
let fenic_type = serde_json::from_str::<FenicType>(&kwargs.dest_dtype)
.expect("Invalid Fenic type string");
Ok(Field::new("casted".into(), fenic_type.canonical_polars_type()))
}
#[polars_expr(output_type_func_with_kwargs=fenic_dtype_str_to_polars_dtype)]
fn cast_expr(inputs: &[Series], kwargs: CastKwargs) -> PolarsResult<Series> {
let src_type = serde_json::from_str::<FenicType>(&kwargs.source_dtype)
.expect("Invalid Fenic type string");
let dest_type = serde_json::from_str::<FenicType>(&kwargs.dest_dtype)
.expect("Invalid Fenic type string");
cast_series_to_fenic_dtype(&inputs[0], &src_type, &dest_type)
}
```
教训:`output_type_func_with_kwargs` 允许从插件的配置(在 schema 解析时解码)计算出其结果类型。这使得将完全外来的类型系统(fenic 的逻辑类型,包括像固定大小嵌入数组这样的东西)接入 Polars 的物理 dtype,同时保持惰性 schema 解析完好无损成为可能。你的插件必须事先、在没有数据的情况下,精确声明它将产生什么形状。
(一个诚实的警告,因为我们把这个作为范例展示。真实代码在格式错误的类型描述符上使用了 `.expect()`,因此一个错误的 JSON 字符串会使 schema 解析路径 panic,而不是提供一个干净的 `PolarsError`。加强它会将其映射为 `ComputeError`。)
如果把这个搞错,声明了 `List` 却构建了 `Float64`,你不会得到编译错误。你会得到一个运行时失败,当产生的列与承诺的 schema 不匹配时。这个声明是一个契约,你需要负责维护它。
### 正确处理并行性和广播
每个 fenic 插件 pas
相似文章
Diplomat:面向 Rust 库的多语言 FFI
Diplomat 是一个多语言单向 FFI 工具,用于封装 Rust 库,旨在将 Rust API 暴露给 C++、JS、Dart 和 JVM 等语言,而无需 FFI 专业知识,填补了 Rust 工具生态系统中的空白。
我们如何(及为何)将生产环境的C++前端基础设施重写为Rust
NearlyFreeSpeech.NET 将其生产环境的C++前端基础设施(nfsncore)重写为Rust,该系统负责所有传入请求的路由、缓存和访问控制。迁移的动机是Rust的安全性保证、性能、生态系统优势以及老化的C++代码库的局限性。
@MaximeRivest: 我用了大约3个提示,让Fable将dplyr移植到Python,以polars/duckdb为后端。我尚未完成对移植的评估,但它……
Maxime Rivest 宣布将 dplyr 移植到 Python,使用 polars/duckdb 作为后端,并评价说移植很扎实。
Rust类型系统中的Lisp
一个嵌入在Rust trait系统中的Lisp解释器,支持在编译时进行递归函数、闭包和延续传递风格。
@blackanger: 从 Epic Lora 里蒸馏了不少好的实践
An experiment on using Rust's type system as an AI coding specification, distilling practices from Epic Lora.