Daft 文本处理实战指南:文本嵌入生成、多 Provider 接入与分块策略
【免费下载链接】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 的文本模态处理能力展开,完整讲解如何用embed_text将文本转换为捕捉语义的数值向量(嵌入),如何在 Sentence Transformers、OpenAI、LM Studio 三类 Provider 之间切换,以及如何用 spaCy 将长文本切分为适合嵌入与检索的小块。读完本文,你将掌握从原始文本到可检索向量的完整链路:安装可选依赖、按场景选择 Provider、处理模型上下文长度约束,以及把嵌入写入向量数据库(如 Turbopuffer)的端到端方案。
背景:为什么在 Daft 中处理文本
Daft 是面向 AI 与多模态负载设计的高性能数据引擎,其模态体系覆盖文本、图片、音频、视频与文档(见 模态总览)。文本处理位于这一体系的基石位置:大语言模型与 RAG 检索的输入输出,本质上都是文本或其向量化表示。
在 Daft 中,文本处理通常以 DataFrame 列(String 类型)为载体,通过表达式(Expression)声明式地完成"读入文本 → 嵌入 → 分块 → 检索/写入"的流水线。这意味着文本 AI 能力可以直接与read_huggingface、read_parquet等任意数据源组合,并天然获得分布式执行的扩展能力。
安装与依赖准备
文本嵌入默认走 Sentence Transformers Provider,需要安装可选依赖(完整可选依赖清单见 安装指南):
pip install -U "daft[transformers]"如果后续要使用 OpenAI Provider(包括 LM Studio 这类 OpenAI 兼容接口),则安装:
pip install -U "daft[openai]"其他常用扩展还有daft[ray](Ray 分布式)、daft[turbopuffer](Turbopuffer 向量库写入)等,可按需组合。若希望一次性安装全部依赖,可执行pip install -U "daft[all]"。
生成文本嵌入:embed_text 基础用法
文本嵌入(Text Embedding)将文本转化为捕捉语义信息的数值向量,是语义搜索、相似度计算及其他 NLP 任务的基础。
最简洁的用法是直接对 DataFrame 中的文本列调用embed_text。默认情况下它使用 Sentence Transformers Provider,无需额外配置:
import daft from daft.functions.ai import embed_text ( daft.read_huggingface("togethercomputer/RedPajama-Data-1T") .with_column("embedding", embed_text(daft.col("text"))) .show() )这里read_huggingface直接从 Hugging Face 加载公开数据集,with_column以声明式方式把"文本嵌入"这一计算挂到新列embedding上。从源码看,embed_text的签名位于 daft/functions/ai/init.py,核心参数包括:
| 参数 | 类型 | 说明 |
|---|---|---|
text | Expression | 输入文本列表达式(必填) |
provider | str \| Provider \| None | 使用的 Provider 名称或实例,缺省时按解析规则回退到默认 Provider |
model | str \| None | 嵌入模型名称,缺省使用 Provider 的默认模型 |
dimensions | int \| None | 输出向量维度(需 Provider 与模型支持),缺省为模型默认维度 |
**options | 任意 | 透传给模型的其他选项,如extra_body、batch_size等 |
执行后新列类型为Embedding[Float32; N],例如使用sentence-transformers/all-MiniLM-L6-v2时输出 384 维向量(该用法示例见 daft/functions/ai/init.py 的 docstring)。
Provider 解析机制
embed_text底层通过_resolve_provider完成 Provider 解析(见 daft/functions/ai/init.py),其优先级如下:
- 显式传入
Provider实例时直接使用; - 传入字符串且当前 Session 中已注册同名 Provider 时,从 Session 加载;
- 传入字符串时调用
load_provider按名称加载; - 使用当前 Session 的 current provider;
- 最终回退到该 API 的默认 Provider(
embed_text默认transformers,prompt默认openai)。
这一机制意味着:即便不显式传provider,只要你在 Session 中配置过 Provider,调用也能自动命中。
使用不同 Provider 生成嵌入
Sentence Transformers(本地开源模型)
Sentence Transformers 是计算嵌入的主流开源模块。安装可选依赖后,用transformersProvider 搭配 Hugging Face 上的任意开源模型即可,例如BAAI/bge-base-en-v1.5:
import daft from daft.functions.ai import embed_text provider = "transformers" model = "BAAI/bge-base-en-v1.5" ( daft.read_huggingface("togethercomputer/RedPajama-Data-1T") .with_column("embedding", embed_text(daft.col("text"), provider=provider, model=model)) .show() )从实现看,transformersProvider 的嵌入器位于 daft/ai/transformers/protocols/text_embedder.py:它通过SentenceTransformer(model_name_or_path, trust_remote_code=True)加载模型,默认批大小batch_size=64,推理在torch.inference_mode()下执行并以convert_to_numpy=True输出;dimensions参数会先经AutoConfig读取模型hidden_size校验(请求维度不得超过模型输出维度),再通过truncate_dim截断。这种"模型维度自动发现 + 显式校验"的设计,避免了用户手填维度与实际模型不符的问题。
OpenAI(云 API)
OpenAI 是文本嵌入的常见云端选择。安装daft[openai]并设置OPENAI_API_KEY环境变量后:
import daft from daft.functions.ai import embed_text provider = "openai" model = "text-embedding-3-small" ( daft.read_huggingface("Open-Orca/OpenOrca") .with_column("embedding", embed_text(daft.col("response"), provider=provider, model=model)) .show() )需要注意模型约束:不同的嵌入模型有不同的上下文长度上限。例如 OpenAI 的text-embedding-3-small最大上下文为 8,192 tokens,超限时会收到类似下面的错误:
openai.BadRequestError: Error code: 400 - {'error': {'message': "This model's maximum context length is 8192 tokens, however you requested 12839 tokens (12839 in your prompt; 0 for the completion). Please reduce your prompt; or completion length.", 'type': 'invalid_request_error', 'param': None, 'code': None}}解决思路有两种:换用上下文更长的模型,或在嵌入前先把文本切分成更小的片段(见下文"分块"章节)。
从源码看,OpenAI Provider 内置了模型档案表(见 daft/ai/openai/protocols/text_embedder.py),记录了官方嵌入模型的元数据:
| 模型 | 默认维度 | 支持覆盖维度 | 输入 token 上限 |
|---|---|---|---|
text-embedding-ada-002 | 1536 | 否 | 8192(默认) |
text-embedding-3-small | 1536 | 是 | 8192(默认) |
text-embedding-3-large | 3072 | 是 | 8192(默认) |
当使用官方base_url时,未登记模型会直接报错;自定义base_url(如接入 OpenAI 兼容网关)时则放行,未知模型的输入 token 上限按 8192 兜底。dimensions只允许在模型支持覆盖维度时传入(text-embedding-3-*系列支持,ada-002不支持)。
LM Studio(本地离线推理)
LM Studio 是本地 AI 模型平台,可在自有机器上运行 Qwen、Mistral、Gemma、gpt-oss 等模型,适合隐私敏感或离线场景。由于 LM Studio 暴露 OpenAI 兼容 API,Daft 直接复用daft[openai]依赖即可:
pip install -U "daft[openai]"LM Studio 默认监听localhost:1234,Daft 中可通过base_url自定义。以下示例使用nomic-ai/nomic-embed-text-v1.5嵌入模型:
import daft from daft.functions import embed_text # 配置 LM Studio Provider(base_url 可选,不传则用 LM Studio 默认地址) daft.set_provider("lm_studio", base_url="http://127.0.0.1:1234/v1") # 选择已在 LM Studio 中加载的文本嵌入模型 model = "text-embedding-nomic-embed-text-v1.5" ( daft.read_huggingface("Open-Orca/OpenOrca") .with_column("embedding", embed_text(daft.col("response"), provider="lm_studio", model=model)) .show() )daft.set_provider用于在当前 Session 内全局配置 Provider(具体用法可参见 AI Providers 总览)。LM Studio 的嵌入器实现位于 daft/ai/lm_studio/protocols/text_embedder.py,有两点设计值得注意:
- 动态维度探测:LM Studio 可加载不同维度的模型,Daft 会先向本地服务发送一次 "dimension probe" 请求读取真实输出维度,避免硬编码;
- 请求参数控制:默认
batch_size=64、max_retries=3,批 token 上限batch_token_limit默认 300,000,并复用 OpenAI 的模型档案做输入 token 上限估算。
面向 OpenAI 兼容服务的扩展参数
对于 OpenAI 兼容的嵌入服务,可通过extra_body透传额外的 JSON 请求属性。例如针对 vLLM 服务,可以让长输入从右侧截断、仅保留前 4,096 个 token(完整示例见 嵌入函数详解):
df = df.with_column( "embedding", embed_text( daft.col("text"), provider=provider, model="BAAI/bge-m3", extra_body={ "truncate_prompt_tokens": 4096, "truncation_side": "right", }, ), )使用嵌入:语义相似度与向量数据库
生成嵌入后最常见的用途是相似度检索或配合向量数据库做 RAG。Daft 内置cosine_distance函数可直接计算向量间余弦距离(数值越小越相似),一个朴素但完整的语义搜索流程如下(示例出自 docs/ai-functions/embed.md):
import daft from daft.functions import embed_text, cosine_distance # 构建知识库 documents = daft.from_pydict({ "doc_id": [1, 2, 3, 4], "text": [ "Python is a high-level programming language", "Machine learning models require training data", "Daft is a distributed dataframe library", "Embeddings capture semantic meaning of text", ] }).with_column( "embedding", embed_text(daft.col("text"), provider="openai", model="text-embedding-3-small") ) # 构造并嵌入查询 query = daft.from_pydict({"query_text": ["What is Daft?"]}).with_column( "query_embedding", embed_text(daft.col("query_text"), provider="openai", model="text-embedding-3-small") ) # 交叉连接后按余弦距离排序取最相似文档 results = ( query.join(documents, how="cross") .with_column("distance", cosine_distance(daft.col("query_embedding"), daft.col("embedding"))) .sort("distance") .select("query_text", "text", "distance") ) results.show()示例输出中,与 "What is Daft?" 距离最近(0.362)的正是 "Daft is a distributed dataframe library",直观体现了嵌入的语义捕捉能力。在大规模场景下,官方建议:预计算并存储嵌入以避免重复生成、采用 ANN 近似最近邻算法、在计算相似度前先过滤文档以减少比较次数、以及将结果写入向量数据库以高效检索。
将嵌入写入向量数据库时,Daft 提供了原生的write_turbopuffer写入器(详见 Turbopuffer 连接器文档),可在一条流水线中完成"读数据 → 嵌入 → 写入向量库"的全过程。
将长文本切分为小块
处理大文本时通常需要先分块(Chunking)。分块粒度取决于场景:句子级分块适用于文档结构不明确或内容多样的数据;段落级分块适合需要跨句保持上下文的 RAG 应用;章节级分块适合结构清晰的超长文档;固定大小切分实现简单但可能切断语义(完整策略讨论见 文本嵌入端到端教程)。
基于 spaCy 的句子分块
spaCy 的句子边界检测比简单的标点切分更稳健,是句子级分块的常用工具。先安装 spaCy 与语言模型:
pip install -U spacy python -m spacy download en_core_web_sm然后在 Daft 中以 用户自定义函数(UDF) 形式封装分块逻辑:
import daft import typing nlp_model_name = "en_core_web_sm" @daft.func def chunk_by_sentences(text: str) -> typing.Iterator[str]: import spacy nlp = spacy.load(nlp_model_name) for sentence in nlp(text).sents: yield sentence.text ( daft.read_huggingface("togethercomputer/RedPajama-Data-1T") .limit(8) .with_column("chunks", chunk_by_sentences(daft.col("text"))) .show() )这里@daft.func将普通 Python 函数提升为可向量化执行、可分布式调度的 UDF:函数按行接收文本并yield多个句子,Daft 自动把每行产出的一组句子组织为 List 列chunks,随后可配合explode将列表展开为独立行、再做嵌入。
生产级流水线:从分块到写入向量库
将上述能力组合起来,即可构建覆盖"分块 → 嵌入 → 写入向量库"的端到端流水线。文本嵌入端到端教程 给出了一份可运行的完整脚本,其核心结构包括:
- 类 UDF 承载模型状态:
ChunkingUDF(加载 spaCy 模型,max_concurrency限制并发实例数)与EncodingUDF(加载 SentenceTransformer 模型到 GPU,使用bfloat16降低显存占用); - Ray 分布式执行:通过
daft.set_runner_ray()将调度切换到 Ray 集群,配合daft.io.IOConfig配置 S3 读取; - 流水线组装:
read_parquet读入 →with_column分块 →explode展开句子 → 生成嵌入 → 构造唯一id→select精简列 →write_turbopuffer写入向量库。
其中write_turbopuffer的核心参数包括namespace(命名空间)、region(区域)、id_column(主键列)、vector_column(向量列)与distance_metric="cosine_distance"(距离度量)。
该教程还给出了调优要点:SENTENCE_TRANSFORMER_BATCH_SIZE越大吞吐越高、越小显存占用越低;NUM_GPU_NODES与CHUNKING_PARALLELISM需按集群规模调整;GPU 显存不足或超出模型最大序列长度时通常需要调小批大小;模型只在每个 worker 加载一次,初始化成本被摊薄;量化(bfloat16/float16)可显著降低显存并提升吞吐。
文本嵌入流水线执行期间多 GPU 接近 100% 利用率的监控图
得益于 Daft 将网络 I/O、CPU 分块与 GPU 推理流水线化,这类流水线在集群上运行时可以获得接近 100% 的 GPU 利用率(如上图所示),从而高效处理数百万级文本文档。
延伸阅读
- 嵌入函数详解(embed_text / embed_image / cosine_distance)
- AI Providers 总览(set_provider、Session 与多 Provider 管理)
- 文本嵌入端到端教程(分块 + 嵌入 + Turbopuffer 写入)
- Turbopuffer 向量数据库连接器
- 自定义函数(UDF)指南
- 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),仅供参考