用 Rust 为 Daft 构建第一个原生扩展:Hello Native Extension 实战指南
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
本指南以 Daft 仓库中的最小原生扩展示例 examples/hello 为核心,完整讲解如何用 Rust 编写一个可被 Daft 分布式执行引擎加载的原生扩展(Native ABI Extension),包括标量函数、聚合函数的实现、Python 包装层、构建安装与测试全流程。读完本文,你将掌握 Daft 扩展机制的底层原理(Arrow C Data Interface ABI、Session 注册模型),并能独立从零搭建一个可 pip 安装、可在 DataFrame 表达式中直接调用的原生扩展包。
从 README 出发:hello 示例解决什么问题
仓库中 examples/hello/README.md 开宗明义:这是一个minimal Daft native extension example(最小 Daft 原生扩展示例),其快速开始只有两步:
# Install the extension in the project .venv uv pip install -e . # Run the tests! pytest -vREADME 同时将读者引导至扩展指南。需要说明的是,指南实际位于 docs/extensions/overview.md(扩展模型总览)与 docs/extensions/authoring.md(Rust 原生扩展编写指南)。根据 overview 的划分,Daft 扩展有两条技术路线:
- Python UDF 型扩展:基于
@daft.func、@daft.cls等自定义代码 API(见 docs/custom-code/func.md),适合编排 Python 生态、外部服务、ML 模型等场景,无需编译原生库; - 原生 ABI 型扩展:以共享库(shared library)为载体,基于 Arrow C Data Interface 的稳定 C ABI,获得底层的向量化原生执行性能。
hello示例正是这条路线的最小实现。
项目骨架:Rust 工程与 Python 包如何组织
从仓库目录结构可以还原出hello扩展的完整布局,它由三部分拼装而成:
examples/hello/ ├── Cargo.toml # Rust 工程配置(cdylib 产物) ├── pyproject.toml # Python 打包配置(setuptools-rust) ├── setup.py # RustExtension 接线文件 ├── hello/ │ ├── __init__.py # Python 表达式包装层 │ └── py.typed # PEP 561 类型标记 ├── src/ │ └── lib.rs # Rust 原生实现(模块 + 函数) └── tests/ └── test_hello.py # pytest 测试Cargo.toml:必须以 cdylib 编译
examples/hello/Cargo.toml 的关键点在于crate-type = ["cdylib"]——只有动态库才能在运行时被dlopen加载进 Daft 进程:
[workspace] [package] name = "hello" edition = "2024" version = "0.1.0" [lib] name = "hello" crate-type = ["cdylib"] [dependencies] daft-ext = {path = "../../src/daft-ext", features = ["arrow-58"]} arrow-array = {version = "58", features = ["chrono-tz"]} arrow-schema = "58"依赖方面有两点值得注意:
daft-ext是官方扩展 SDK,位于 src/daft-ext,提供类型、trait 与宏;这里以仓库相对路径path方式引用,正式发布时应改为版本号依赖;- feature flag
arrow-58必须与你的 arrow-rs 版本匹配。ABI 边界使用纯 C 结构体(ArrowSchema、ArrowArray),因此扩展并不绑定 Daft 内部的 arrow-rs 版本;在daft-ext上启用与自身 arrow-rs 一致的特性(arrow-56、arrow-57或arrow-58)即可获得安全的.into()转换,不支持的版本则可使用from_owned/into_owned/from_raw/as_raw逃生舱接口。
pyproject.toml:setuptools-rust 构建系统
examples/hello/pyproject.toml 使用setuptools+setuptools-rust作为构建后端,把 Rust cdylib 的编译集成进 Python 打包流程:
[build-system] requires = ["setuptools", "setuptools-rust"] build-backend = "setuptools.build_meta" [project] name = "hello" version = "0.1.0" requires-python = ">=3.10" dependencies = ["daft"] [project.optional-dependencies] test = ["pytest"]包本身依赖daft运行时;在仓库内开发时通过[tool.uv.sources]将其指向仓库根目录的可编辑安装:
[tool.uv.sources] daft = {path = "../..", editable = true}setup.py:把编译产物装进 Python 包目录
examples/hello/setup.py 是 Python 与 Rust 的接线文件:
from setuptools import find_packages, setup from setuptools_rust import Binding, RustExtension setup( packages=find_packages(), rust_extensions=[ RustExtension( "hello.libhello", path="Cargo.toml", binding=Binding.NoBinding, strip=True, ) ], )RustExtension("hello.libhello", ...)的含义是:把编译出的libhello.so放进hello/包目录内,这样Session.load_extension才能定位到它;binding=Binding.NoBinding是因为 Daft 扩展导出的是裸 C 符号(daft_module_magic),而非 PyO3 绑定;strip=True用于裁剪符号表减小体积。
扩展入口:#[daft_extension]模块与 install 钩子
原生扩展的 Rust 侧入口在 examples/hello/src/lib.rs 中。一个扩展由两部分组成:模块(module,入口点)与一个或多个函数(标量或聚合)。
use std::{ffi::CStr, sync::Arc}; use arrow_array::{Array, ArrayRef}; use arrow_schema::{DataType, Field}; use daft_ext::{daft_extension, prelude::*}; // ── Module ────────────────────────────────────────────────────────── #[daft_extension] struct HelloExtension; impl DaftExtension for HelloExtension { fn install(session: &mut dyn DaftSession) { session.define_function(Arc::new(Greet)); session.define_aggregate_function(Arc::new(StringCount)); } }机制上(详见 docs/extensions/authoring.md):
#[daft_extension]宏会生成daft_module_magic这个 C 符号,Daft 运行时dlopen加载共享库时就是靠它发现扩展入口的;impl DaftExtension中的install(session: &mut dyn DaftSession)是扩展安装钩子,在扩展被加载进某个会话时调用一次,所有函数(标量用define_function、聚合用define_aggregate_function)都在这里注册;- 函数名在会话内是全局的,authoring 指南建议使用前缀(如
myext_greet)避免多扩展共存时冲突。
标量函数:greet 的 Rust 实现
当前仓库的hello示例使用#[daft_func]宏实现标量函数,写法非常简洁:
// ── Scalar Function ──────────────────────────────────────────────── #[daft_func] fn greet(name: &str) -> String { format!("Hello, {}!", name) }authoring 指南则展示了更底层的等价写法:每个标量函数是一个实现DaftScalarFunctiontrait 的结构体,包含三个方法:
name(&self) -> &CStr:函数名,Python 侧通过它查找(注意使用c"..."字符串字面量);return_field(&self, args: &[ArrowSchema]) -> DaftResult<ArrowSchema>:类型检查与输出类型声明。输入字段以 C Data Interface 的ArrowSchema形式传入,可用.as_raw()零拷贝借用为 arrow-rs 的FFI_ArrowSchema做校验(如检查参数个数、是否字符串类型),再.into()转换回 ABI 类型返回;call(&self, args: Vec<ArrowData>) -> DaftResult<ArrowData>:真正的执行逻辑,接收整列ArrowData,输出整列结果——所有数据都走 Arrow 数组,没有任何逐行 Python 开销。
fn call(&self, args: Vec<ArrowData>) -> DaftResult<ArrowData> { // 1. 取出输入列并转换为 arrow-rs FFI 类型 // 2. 用 StringBuilder 逐行构建输出(null 也要保留) // 3. 通过 arrow::ffi::to_ffi 转回 ABI 类型返回 }实现中有两个容易踩坑的点(authoring 指南明确提醒):
- 字符串类型必须是
LargeUtf8(i64 偏移):Daft 内部字符串使用 i64 offsets,向下转型时必须用as_string::<i64>(),用i32会在运行时 panic;return_field做类型校验时也要接受DataType::LargeUtf8; - 错误分类:
return_field中的 schema 违规返回Err(DaftError::TypeError(...)),call中的执行失败返回Err(DaftError::RuntimeError(...))。
聚合函数:string_count 的三阶段管线
除了标量函数,扩展还可以注册聚合函数(UDAF)。仓库中的StringCount(统计非空字符串个数)是一个完整的实现范例(examples/hello/src/lib.rs)。聚合函数遵循三阶段管线:
- 聚合(
aggregate):把输入数组加工成部分状态(partial state); - 合并(
combine):把多个部分状态合并成一个; - 终结(
finalize):从合并后的状态产出最终标量结果。
状态以Vec<ArrowData>交换,每个元素对应一个状态字段,FFI 层会把这些字段透明地打包成一个 Struct 数组。对应 trait 的实现要点:
struct StringCount; impl DaftAggregateFunction for StringCount { fn name(&self) -> &CStr { c"string_count" } // 输出类型声明:1 个参数,否则 TypeError fn return_field(&self, args: &[ArrowSchema]) -> DaftResult<ArrowSchema> { // ... 校验 args.len() == 1,返回 Int64 字段 "string_count" } // 中间状态字段声明:单个 Int64 字段 "count" fn state_fields(&self, _args: &[ArrowSchema]) -> DaftResult<Vec<ArrowSchema>> { // ... } // 阶段一:统计输入列的非空个数,产出单行状态 fn aggregate(&self, inputs: Vec<ArrowData>) -> DaftResult<Vec<ArrowData>> { let non_null_count = input.len() - input.null_count(); // 返回 Int64Array 状态 } // 阶段二:把多行状态求和 fn combine(&self, states: Vec<ArrowData>) -> DaftResult<Vec<ArrowData>> { // 遍历 counts 数组求和(跳过 null) } // 阶段三:取出最终 count fn finalize(&self, states: Vec<ArrowData>) -> DaftResult<ArrowData> { // ... } }register 时使用session.define_aggregate_function(Arc::new(StringCount)),与标量函数并存于同一个install钩子中。这种三阶段设计天然适配分布式执行:每个分区先在本地aggregate出部分状态,再在 shuffle 后combine,最后finalize产出结果。
Python 包装层:把原生函数接入 Expression DSL
原生函数需要在 Python 侧提供符号,才能在 DataFrame 表达式 DSL 中使用。包装层在 examples/hello/hello/init.py:
from __future__ import annotations from typing import TYPE_CHECKING import daft if TYPE_CHECKING: from daft.expressions import Expression def greet(name: Expression) -> Expression: """Greet someone by name.""" return daft.get_function("greet", name) def string_count(name: Expression) -> Expression: """Count non-null strings.""" return daft.get_aggregate_function("string_count", name)- 标量函数通过
daft.get_function("greet", name)解析——它调用当前会话的get_function,按DaftScalarFunction::name()注册的名字查找函数并绑定参数(模块级 API 见 daft/session.py); - 聚合函数通过
daft.get_aggregate_function(...)走同样的解析逻辑(对应DaftAggregateFunction::name()); - 加一层 Python 包装的意义在于:可以用类型注解(
Expression -> Expression)、docstring 提供自动补全与文档,还可以在链接到原生符号之前做参数预处理。authoring 指南特别指出,SQL 解析函数并不需要这层 Python(所以扩展不依赖 PyO3),但 Python 函数让 Expression DSL 用起来更符合 Python 习惯。
在 Python 包内放置空的py.typed标记文件(本示例即如此),即可获得类型检查器的完整支持。
构建、安装与测试:从源码到可运行
安装
按照 README 的快速开始,在项目虚拟环境中执行:
uv pip install -e .setuptools-rust会先编译 Rust cdylib,再以可编辑模式安装 Python 包;开发时通过[tool.uv.sources]将daft指向仓库根目录,保证扩展依赖的 Daft 与当前源码一致。
加载与使用
import daft # Step 1. Import your extension module import hello # Step 2. Load the extension into the current daft session daft.load_extension(hello) # Step 3. Use in your dataframe! df = daft.from_pydict({"name": ["John", "Paul"]}) df = df.select(hello.greet(df["name"])) df.show()输出:
╭──────────────╮ │ greet │ │ --- │ │ String │ ╞══════════════╡ │ Hello, John! │ ├╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤ │ Hello, Paul! │ ╰──────────────╯daft.load_extension的模块级入口定义在 daft/session.py,它接受字符串、模块对象或路径,最终委托给会话的Session.load_extension(daft/session.py)。
测试验证
examples/hello/tests/test_hello.py 提供了完整的 pytest 覆盖,直接pytest -v即可运行:
test_greet/test_greet_null:验证标量函数对普通值与None的处理(null 必须原样透传);test_greet_show:验证.show()的可视化输出;test_string_count/test_string_count_with_nulls:验证聚合函数对含 null 数据的非空计数;test_aggregate_not_available_without_extension:验证未加载扩展的会话中调用string_count会抛异常——这从反面印证了函数是按会话隔离注册的。
测试里使用了独立作用域会话而非全局活动会话:
sess = Session() sess.load_extension(hello) with sess: result = df.select(greet(col("name"))).collect().to_pydict()会话作用域与最佳实践
作者指南对会话机制有几个关键说明,直接关系到扩展的正确使用:
- 扩展在进程内只加载一次:重复调用
load_extension对同一共享库只会dlopen一次,会话只是名字解析的作用域机制; - 函数仅在加载了该扩展的会话中可用:
with sess:上下文管理器可以把查询限定到特定会话,避免污染全局会话; - 命名前缀:函数名在会话内全局唯一,定义多个函数或可能与其他扩展共存时,建议使用
<extension_name>_<fn_name>前缀防止冲突; - 类型检查器:Python 包装函数统一使用
TYPE_CHECKING守卫导入Expression,并给每个函数加上类型注解与 docstring。
延伸:更多原生扩展范式
hello是理解 Daft 原生扩展机制的最佳起点,仓库中还有两个同构但更具实战价值的示例可以继续研读:
- examples/dvector:pgvector 风格的原生扩展,提供
l2_distance、cosine_distance、inner_product、jaccard_distance等向量距离函数,展示了多个函数并存注册的组织方式; - examples/hello_cpp:纯 C++ 原生扩展,基于 Apache Arrow C++ 直接使用 Daft 的裸 C ABI,说明只要语言能产出共享库、导出约定的 C ABI 并能读写 Arrow C Data Interface 数组,就可以接入 Daft。
此外,扩展 SDK 的完整定义位于 src/daft-ext/src(含abi/、ffi/、function.rs、aggregate.rs、session.rs等模块),阅读它可以进一步理解 ABI 结构体与 trait 的底层契约。对于想直接用 Python 快速构建领域扩展(AI 推理、文件处理等)的场景,则可以参考 docs/extensions/overview.md 中介绍的 UDF 型扩展路线。
【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考