news 2026/9/18 11:21:35

用 Rust 为 Daft 构建第一个原生扩展:Hello Native Extension 实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
用 Rust 为 Daft 构建第一个原生扩展:Hello Native Extension 实战指南

用 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 -v

README 同时将读者引导至扩展指南。需要说明的是,指南实际位于 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 flagarrow-58必须与你的 arrow-rs 版本匹配。ABI 边界使用纯 C 结构体(ArrowSchemaArrowArray),因此扩展并不绑定 Daft 内部的 arrow-rs 版本;在daft-ext上启用与自身 arrow-rs 一致的特性(arrow-56arrow-57arrow-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 的结构体,包含三个方法:

  1. name(&self) -> &CStr:函数名,Python 侧通过它查找(注意使用c"..."字符串字面量);
  2. return_field(&self, args: &[ArrowSchema]) -> DaftResult<ArrowSchema>:类型检查与输出类型声明。输入字段以 C Data Interface 的ArrowSchema形式传入,可用.as_raw()零拷贝借用为 arrow-rs 的FFI_ArrowSchema做校验(如检查参数个数、是否字符串类型),再.into()转换回 ABI 类型返回;
  3. 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)。聚合函数遵循三阶段管线:

  1. 聚合(aggregate:把输入数组加工成部分状态(partial state);
  2. 合并(combine:把多个部分状态合并成一个;
  3. 终结(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_distancecosine_distanceinner_productjaccard_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.rsaggregate.rssession.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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/18 11:20:28

数据库系统原理实战:从关系模型到ACID落地

1. 这不是教科书笔记&#xff0c;而是一份“能跑通、能排错、能讲清楚”的数据库系统原理实战手记我带过三届数据库课程设计&#xff0c;也给金融、制造、政务类客户做过十多个数据库架构优化项目。每次新人一上来就翻《数据库系统概论》第六版&#xff0c;划重点、背定义、抄E…

作者头像 李华
网站建设 2026/9/18 11:20:20

洪水调节课程设计:从水量平衡到水库调洪演算全流程解析

简介&#xff1a;一份面向水利水电工程专业学生的洪水调节课程设计参考文档&#xff0c;完整展示三峡大学该课程设计的任务要求与计算思路。内容涵盖设计目的、工程基本资料、洪水标准确定&#xff0c;以及列表试算法、半图解法推求下泄流量、库容与水位变化过程的详细流程&…

作者头像 李华
网站建设 2026/9/18 11:20:02

MindSpore Transformers训练实时监控实战:Callback+WebSocket+ECharts

前阵子我调一个 Deformable DETR 的微调实验&#xff0c;模型用 MindSpore Transformers 套件加载&#xff0c;睡前看了一眼 loss 还在 0.8 附近&#xff0c;心里想着“还行&#xff0c;明早起来应该能收敛”。结果第二天打开终端&#xff0c;屏幕上一行刺眼的 loss: nan&#…

作者头像 李华
网站建设 2026/9/18 11:19:40

Highcharts表格直驱可视化:HTML Table自动转图表教程

1. 为什么这个标题值得你花15分钟认真读完Highcharts 实战&#xff5c;HTML表格数据源自动可视化开发教程&#xff08;附Demo&#xff09;——这行字里藏着三个关键信号&#xff1a;Highcharts是工业级图表库的“老司机”&#xff0c;不是玩具级轮子&#xff1b;HTML表格数据源…

作者头像 李华
网站建设 2026/9/18 11:17:58

【ComfyUI】QwenImageEdit 基础图生图

今天展示的案例是一个基于 Qwen-Image 编辑功能的 ComfyUI 工作流。该工作流围绕图像编辑展开,通过加载扩散模型、文本编码器、VAE 模块以及 LoRA 适配器,结合输入图像与文本提示,实现对图像中元素的精准移除与增强效果。 整体设计不仅保证了画面质量,也通过采样器与归一化…

作者头像 李华