资讯动态

Daft 扩展机制完全指南:Python UDF 与原生 ABI 双路径构建领域专属能力

发布时间:2026/9/17 19:12:01 来源:尧图企业网站定制
Daft 扩展机制完全指南Python UDF 与原生 ABI 双路径构建领域专属能力【免费下载链接】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/DaftDaft 扩展Extensions是可复用的库用于在 Daft 之上叠加领域专属能力——无论是地理空间函数、向量检索、HTML 文档处理还是面向特定数据库的读写连接器都不必把每一种特化函数、数据类型、集成或工作流塞进 Daft 核心。本文基于 Daft 扩展总览 展开梳理两种扩展构建路径Python UDF 扩展与原生 ABI 扩展、它们在仓库中的真实落地案例daft.functions.ai、文件 API、daft-lance、hello/dvector/hello_cpp并深入源码剖析扩展的加载与注册原理。读完你将掌握如何判断选型、如何编写并打包一个 Daft 扩展、如何加载进 Session 并在 DataFrame 表达式中调用。扩展机制总览为什么 Daft 需要扩展Daft 的核心是一个高性能数据引擎面向 AI 与多模态负载图像、音频、视频、结构化数据。扩展机制的存在是为了在不污染核心代码的前提下让社区与领域团队以「可复用的库」形式扩展 Daft 的能力边界。构建扩展有两条宽泛的路径路径技术底座典型能力是否依赖原生共享库Python UDF 扩展Daft 的 custom-code APIdaft.func、daft.func.batch、daft.cls、daft.method.batch编排既有 Python 生态、外部服务、ML 模型、GPU、数据维护任务否普通 Python 包即可原生 ABI 扩展基于 Arrow C Data Interface 的稳定 C ABI高性能标量函数、聚合函数UDAF、Python 表达式包装、扩展型数据类型是随 pip 包分发共享库两条路径有一个共同的用户体验目标最终都向用户暴露「看起来像普通 Daft 表达式函数」的干净 Python API。实现细节被隐藏在函数、类与表达式之后——这正是 Daft 扩展设计哲学的核心。路径一基于 Python UDF 的扩展核心装饰器与适用场景最快构建 Daft 扩展的方式往往是纯 Python。使用daft.func、daft.func.batch、daft.cls和daft.method.batch四个装饰器贡献者即可把可复用的 Python 逻辑打包并在 Daft 分布式执行引擎内部运行。这类扩展不需要原生共享库也不需要 Arrow C ABI——它们是普通的 Python 包只是对外暴露了更高层的 Daft API。最适合的场景包括编排既有 Python 生态与外部服务如调用 HTTP API、数据库客户端运行 ML 模型与 GPU 推理配合 GPU 资源提示数据维护任务压缩、索引、合并等领域专属工作流与多模态处理流水线。daft.func家族的函数变体详见 stateless UDF 文档提供了从行级到批处理的完整性能梯度import daft # 行级默认逐行处理返回每行一个值 daft.func def add_and_format(a: int, b: int) - str: return fSum: {a b} # 异步行级并发执行适合 I/O 密集配合 max_concurrency 限流 daft.func(max_concurrency10) async def fetch_url(url: str) - str: ... # 生成器一行输入产出多行输出 daft.func def split_text(text: str) - Iterator[str]: for word in text.split(): yield word # 批处理整批数据以 daft.Series 传入可对接 PyArrow / NumPy 做向量化计算 daft.func.batch(return_dtypedaft.DataType.int64()) def add_series(a: Series, b: Series) - Series: import pyarrow.compute as pc return pc.add(a.to_arrow(), b.to_arrow())而对于有状态的负载最典型的是 GPU 模型推理应使用daft.cls模型在__init__中只初始化一次跨行复用昂贵的加载开销被摊销到整个查询中用gpus参数申请 GPU用daft.method.batch(batch_size...)控制批大小。仓库内置的 AI 函数见下文正是建立在这套机制之上。典型案例daft.functions.ai —— 把 UDF 藏在表达式后面Daft 自身就是这种模式的最佳实践。daft.functions.ai模块API 见 docs/api/ai.md使用指南见 AI 函数对外暴露高层函数promptprompt 文档embed_text/embed_imageembed 文档classify_text/classify_imageclassify 文档对用户而言这些就是普通的 Daft 表达式函数可以直接放进df.select(...)df df.select( daft.functions.ai.embed_text(df[text]), daft.functions.ai.classify_image(df[image_url], labels[cat, dog]), )而在底层它们使用的是 Daft 的 UDF 与 class-UDF 机制包括批处理batching整批喂给模型以提升 GPU/CPU 吞吐并发控制cpus、gpus、max_concurrency决定并行度而不是放置位置重试retriesmax_retries配合指数退避容忍外部服务抖动GPU 资源提示通过 GPU 使用指南 中的daft.cls(gpus1)模式申请设备。各模型提供商的接入实现位于 daft/ai/ 目录openai、google、transformers、vllm、lm_studio 等子模块与 daft/functions/ai/ 的表达式层分离——这本身就是「实现细节隐藏在表达式 API 之后」的架构示范。典型案例文件 API —— 表达式级的领域积木文件 API 遵循同样的产品模式。daft.File类型值以及file、audio_file、video_file、file_path、file_size等辅助函数源码见 daft/file/ 与 daft/functions/file_.py使用说明见 文件模态为用户提供处理文件的表达式级积木。这些文件对象可以继续流入 Python UDF、模型流水线和领域库而用户完全不需要关心执行细节——例如在 UDF 中直接消费daft.File指向的字节数据做音频转写或图像解码import daft from daft import col daft.func def summarize_audio(file: daft.File) - str: # file 已封装好读取逻辑这里聚焦领域逻辑 return transcribe(file.path) df daft.read_audio(s3://bucket/audio/*.mp3) df df.select(summarize_audio(col(audio)))典型案例daft-lance —— 纯 Python 分布式任务库daft-lance是 Python UDF 扩展模型的社区实例它为 Daft 增加 Lance 专属的分布式维护与数据管理操作——文件压缩compaction、标量索引scalar indexing、列合并column merging与 REST catalog 操作。内部使用 Daft 的 Python UDF 与 class-UDF API把 Lance 任务分布到一次 Daft 查询中用户侧接触到的只是compact_files、create_scalar_index、merge_columns_df这样的普通 Python 函数from daft_lance import compact_files, create_scalar_index compact_files(s3://bucket/my_dataset) create_scalar_index(s3://bucket/my_dataset, columnname, index_typeINVERTED)路径二原生 ABI 扩展Arrow C Data Interface与 Arrow 版本彻底解耦对于需要底层、向量化性能的贡献者Daft 还支持通过基于Arrow C Data Interface的 C ABI 构建原生扩展。这一设计的核心卖点是扩展不耦合任何特定 Arrow 库版本ABI 边界只使用纯 C 结构体ArrowSchema、ArrowArray因此扩展可以使用任意版本的 arrow-rs甚至完全不同的 Arrow 实现扩展以pip 可安装的 Python 包形态分发内部打包一个原生共享库。用户的使用流程固定为三步import扩展模块 →daft.load_extension(扩展模块)加载进 Session → 在 DataFrame 表达式中调用普通 Python 包装函数。原生 ABI 扩展可以添加的能力包括高性能标量函数、聚合函数UDAF、Python 表达式包装以及扩展型逻辑数据类型如DataType.extension(...)见 all_datatypes 文档。仓库内的三个官方示例当前仓库提供了三个可直接研读的完整示例示例语言功能关键文件examples/helloRust最小原生扩展greet标量函数 string_count聚合函数src/lib.rs、Cargo.toml、hello/init.pyexamples/dvectorRustpgvector 风格的向量距离函数l2_distance、inner_product、cosine_distance、l1_distance、hamming_distance、jaccard_distancesrc/lib.rsexamples/hello_cppC使用 Apache Arrow C 与原始 Daft C ABI 的纯 C 扩展src/Rust 目前通过daft-ext拥有最符合人体工学的 SDKsrc/daft-ext/src/lib.rs 是 SDK 的源码入口C 通过原始 ABI 演示。其他系统语言只要能做到三件事同样可行产出共享库、导出约定的 C ABI、能读写 Arrow C Data Interface 数组。原生扩展的工作原理从#[daft_extension]到daft.get_function以examples/hello为例一个原生扩展由两部分组成模块入口点与一个或多个函数标量或聚合。Rust 侧的核心结构见 examples/hello/src/lib.rsuse std::sync::Arc; use daft_ext::{daft_extension, prelude::*}; // #[daft_extension] 会生成 Daft 运行时在 dlopen 加载共享库时寻找的 // daft_module_magic C 符号并把 HelloExtension 转成 hello_extension 作为模块名。 #[daft_extension] struct HelloExtension; impl DaftExtension for HelloExtension { // 扩展安装钩子扩展被加载进 session 时调用一次在此注册每个函数。 fn install(session: mut dyn DaftSession) { session.define_function(Arc::new(Greet)); session.define_aggregate_function(Arc::new(StringCount)); } }注意仓库中的实际实现使用#[daft_func]宏便捷定义greet而聚合函数StringCount则完整实现DaftAggregateFunctiontrait 的四个方法return_field声明输出类型、state_fields声明中间状态字段、aggregate输入数组 → 部分状态、combine合并部分状态、finalize产出最终结果——这是与 Python 侧daft.udaf三阶段管线aggregate / combine / finalize对应的原生版。Python 侧examples/hello/hello/init.py则通过daft.get_function与daft.get_aggregate_function把名字解析到会话中已注册的函数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)这些 Python 包装并非必需SQL 解析函数不依赖 PyO3但为用户提供了类型提示、自动补全与 docstring 的「Pythonic」使用体验。加载机制源码剖析Session 与 dlopen扩展的加载入口在 daft/session.py 的Session.load_extension它接受字符串路径、Path或模块对象三种输入def load_extension(self, extension: str | types.ModuleType | Path) - None: # 1. 解析出共享库路径模块对象则通过 _get_shared_lib 找到 .so # 2. 全局加载共享库使符号对其他库可见 ctypes.CDLL(path, modectypes.RTLD_GLOBAL) # 3. 在 Rust 侧把扩展注册进当前会话 self._session.load_extension(path)几个值得注意的语义从源码与示例测试可确认模块级快捷方式daft.load_extension(...)、daft.get_function(name, *args)、daft.get_aggregate_function(...)会代理到当前活跃 session见 daft/session.py 与 daft/session.pySession 隔离扩展在每个进程内只dlopen一次Session 只是名称解析的作用域函数只在加载了该扩展的 session 中可用。示例测试 examples/hello/tests/test_hello.py 演示了用with sess:上下文管理器把查询限定到指定 session且验证「未加载扩展的 session 调用string_count会抛异常」名称全局性函数名在 session 内全局唯一扩展定义多个函数或被多个扩展共同加载时建议用myext_greet式前缀避免冲突实验性警告load_extension当前是实验性 API调用时会发出警告未来可能变更。SDK 侧src/daft-ext/src/lib.rs 将 ABI 层组织为abi、aggregate、error、function、prelude、session与ffi等模块并通过arrow-56/arrow-57/arrow-58/arrow-59特性开关提供into()/as_raw()等安全转换助手——示例 examples/hello/Cargo.toml 即使用daft-ext { path ../../src/daft-ext, features [arrow-58] }。Crate 必须编译为cdylib以便运行时dlopen并用setuptools-rustsetup.py 中RustExtension(hello.libhello, bindingBinding.NoBinding, stripTrue)把编译产物放入 Python 包目录。扩展可以构建什么综合 docs/extensions/overview.md 与社区实践贡献者可以构建纯 Python UDF 扩展库基于daft.func、daft.func.batch、daft.cls、daft.method.batch分布式任务库如daft-lance那样把 Lance/其他存储的维护操作分发到 Daft 查询GPU 推理扩展用daft.cls(gpus...)包装模型跨行复用设备驻留的模型实例AI / 文件 / 媒体 / 多模态处理库把 UDF 藏在表达式 API 之后daft.functions.ai与文件 API 即仓库内范例原生标量函数以 Arrow C Data Interface 为边界的向量化标量 UDF原生聚合函数 / UDAFaggregate → combine → finalize 三阶段分布式聚合Python 表达式包装提供类型提示与文档的高层 Python 函数扩展型逻辑数据类型如DataType.extension(...)参考 all_datatypes 文档 与daft-geo社区示例中Point2D/Point3D的用法。从哪里开始想直接安装使用现成扩展浏览 社区扩展列表daft-h3、daft-lance、daft-html、daft-geo、daft-qdrant、daft-doris覆盖原生 ABI、Python UDF、自定义 DataSource/DataSink 等不同形态想动手写原生扩展跟随 原生扩展编写指南 的完整教程项目脚手架、Cargo.toml/pyproject.toml/setup.py配置、标量与聚合函数实现、pytest验证并对照仓库内的 examples/hello、examples/dvector、examples/hello_cpp 三个示例想看 Daft 之上的真实项目生态浏览 Built on DaftDeltaCat、SMiR、hypergraph 等以 Daft 为引擎或计算后端想先掌握 Python UDF 基础系统阅读 stateless UDF 文档、stateful class UDF 文档 与 UDAF 文档。选型建议如果你的扩展主要编排 Python 生态、调用外部服务或驱动模型推理纯 Python UDF 路径路径一是最快且最灵活的选择如果你需要每行数据上的极致向量化性能、希望零 Python 逐行开销或者要定义扩展型数据类型则应走原生 ABI 路径路径二。两条路径产出的最终用户体验是一致的——一组像内置函数一样自然的 Daft 表达式。【免费下载链接】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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价