ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

使用 Daft 将 DataFrame 写入 Turbopuffer:向量与全文检索的混合存储实战

使用 Daft 将 DataFrame 写入 Turbopuffer:向量与全文检索的混合存储实战 使用 Daft 将 DataFrame 写入 Turbopuffer向量与全文检索的混合存储实战【免费下载链接】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 将 Turbopuffer 作为官方推荐的搜索后端支持把任意 DataFrame包括图像、音频、视频和结构化数据直接写入 Turbopuffer 命名空间从而同时获得向量相似检索与全文检索能力。本文以 docs/connectors/turbopuffer.md 为核心结合 turbopuffer 数据写入器源码 与 DataFrame.write_turbopuffer 实现完整讲解从数据读取、向量化到写入 Turbopuffer 的端到端流程以及每个写入参数的底层行为读完即可在自己的 RAG、语义搜索或多模态检索场景中直接落地。为什么选择 TurbopufferTurbopuffer 是一个基于对象存储构建的快速搜索引擎其核心特性是将向量搜索与全文搜索full-text search合并在同一份数据上所有数据天然可被检索。相比传统向量数据库它的存储成本更低、写入吞吐更高适合与 Daft 这类面向 AI 与多模态数据的高性能数据处理引擎配合使用Daft 负责大规模数据的读取、清洗、分块与向量化Turbopuffer 负责最终的可检索存储。快速上手从 Hugging Face 数据集写入 Turbopuffer官方示例演示了最典型的应用链路读取 Hugging Face 上的Open-Orca/OpenOrca数据集对response列计算文本向量然后写入 Turbopuffer。安装依赖需要安装带turbopuffer扩展的 Daft以及至少一种向量化后端# 安装 Turbopuffer 写入支持 $ pip install -U daft[turbopuffer] # 二选一OpenAI 托管模型 $ pip install -U daft[openai] # 或者本地开源模型Sentence Transformers $ pip install -U daft[transformers]turbopuffer扩展在 docs/install.md 的安装向导中同样被列为可选组件。另外需要准备两个密钥环境变量TURBOPUFFER_API_KEYTurbopuffer 平台 API 密钥用于身份认证OPENAI_API_KEY仅在选用 OpenAI 向量模型时需要。完整示例代码import os import daft from daft.functions.ai import embed_text # 创建 embedding二选一OpenAI 或 Sentence Transformers # 需要设置 OPENAI_API_KEY if os.getenv(OPENAI_API_KEY): provider openai model text-embedding-3-small else: provider sentence_transformers model BAAI/bge-base-en-v1.5 turbopuffer_config { namespace: daft-tpuf-example, region: gcp-us-central1, distance_metric: cosine_distance, schema: { text: { type: string, full_text_search: True, } } } ( daft.read_huggingface(Open-Orca/OpenOrca) .limit(8) .with_column(vector, embed_text(daft.col(response), providerprovider, modelmodel)) .write_turbopuffer(**turbopuffer_config) # 需要设置 TURBOPUFFER_API_KEY )这段代码做了三件事读取数据daft.read_huggingface(Open-Orca/OpenOrca)直接从 Hugging Face 加载数据集limit(8)限制行数便于快速验证计算向量embed_text根据环境变量自动切换向量模型——存在OPENAI_API_KEY时用 OpenAI 的text-embedding-3-small否则回退到 Hugging Face 上的开源模型BAAI/bge-base-en-v1.5通过 Sentence Transformers 本地推理生成的向量写入vector列写入.write_turbopuffer(**turbopuffer_config)将整个 DataFrame 作为文档写入daft-tpuf-example命名空间。embed_text等 AI 函数的更多用法可参考 daft/functions/ai 与 AI 函数文档更完整的向量化 写入流水线含文本分块、多 GPU 分布式执行见 生成文本向量写入 Turbopuffer 的完整指南。write_turbopuffer 参数详解与底层行为DataFrame.write_turbopuffer是连接 Daft 与 Turbopuffer 的唯一入口其实现位于 daft/dataframe/dataframe.py#L2550-L2597内部创建 TurbopufferDataSink 后经write_sink提交到执行引擎。全部参数如下参数类型说明默认行为namespacestr/Expression目标命名空间必填传表达式时支持多命名空间动态路由api_keystrAPI 密钥为空时读取环境变量TURBOPUFFER_API_KEYregionstrTurbopuffer 区域为空时使用客户端默认区域distance_metriccosine_distance/euclidean_squared向量相似度度量可选写入时透传给namespace.write()schemadict手动指定字段 schema可选透传给namespace.write()id_columnstr用作文档 ID 的列名自动重命名为id缺省时要求数据已有id列vector_columnstr用作向量索引的列名自动重命名为vector缺省时要求数据已有vector列client_kwargsdict透传给turbopuffer.Turbopuffer构造器的参数与显式api_key/region合并write_kwargsdict透传给namespace.write()的参数与显式distance_metric/schema合并行到文档的映射规则从 TurbopufferDataSink 的文档字符串 可以看到写入器把 DataFrame 的每一行转换成一个 Turbopuffer 文档映射规则如下id列是硬性要求每个文档必须有唯一 ID。如果表中没有id列必须通过id_column指定哪一列作为 ID该列在写入时会自动重命名为idvector列按需提供只有当目标命名空间配置了向量索引时才必须存在向量列可通过vector_column指定写入时自动重命名为vector其余所有列自动成为文档属性attributes无需额外声明即可随文档写入配合schema参数中的full_text_search: True即可开启全文检索。_prepare_arrow_table方法turbopuffer_data_sink.py#L140-L155还实现了一个细节ID 为空的行会被自动过滤掉通过 Arrow 的is_null取反过滤避免写入无效文档。命名空间校验规则当namespace传字符串时构造函数会调用_check_namespace_name进行校验turbopuffer_data_sink.py#L41-L52规则包括必须是字符串长度在 1 到 128 个字符之间只能包含字母数字字符以及-、_、.三种符号。不满足任一条件都会抛出带明确提示的ValueError避免把非法命名空间发给远端。多命名空间写入namespace 传表达式namespace参数除了字符串还可以传一个Daft 表达式Expression用于按数据内容动态路由到多个命名空间。例如按category列的值把不同类别的数据写入各自的命名空间。从 write 方法实现 可以看到其内部流程用micropartition.partition_by_value(self._namespace)按表达式计算结果对数据分片对每个分片依次取出其命名空间值逐个做名称合法性校验为每个命名空间创建turbopuffer.Namespace并分别写入。这样一条流水线即可完成数据的水平切分与多库分发无需多次调用。参数的合并与冲突处理client_kwargs与write_kwargs提供了向 Turbopuffer 客户端透传任意参数的通道如重试次数、超时、Turbopuffer 写入端的其他高级选项。需要注意的合并规则turbopuffer_data_sink.py#L111-L130显式传入的api_key、region会合并进client_kwargs显式传入的distance_metric、schema会合并进write_kwargs若同一键在kwargs中已被显式设置再次传入同名参数会抛出ValueError防止静默覆盖。写入结果的错误处理与返回结构暂态错误优雅降级非暂态错误快速失败网络写入难免遇到瞬时故障TurbopufferDataSink 的_write_with_error_handling实现了两档错误策略并有 tests/io/test_turbopuffer_write.py 中的 Mock 测试逐类验证暂态错误gracefulAPIConnectionError、InternalServerError、RateLimitError、ConflictError、APITimeoutError这类错误——Turbopuffer 官方客户端本身已内置重试若重试后仍失败写入器将其记为status: failed的结果并继续后续写入不会中断整个流水线非暂态错误fail fastAuthenticationError401、BadRequestError400、PermissionDeniedError403、NotFoundError404、UnprocessableEntityError422等错误说明请求本身有问题重试无意义异常会直接向上抛出让脚本立即失败以便尽快暴露配置错误。对应测试用例test_resilience_to_transient_errors断言暂态错误下全部行被标记为rows_not_written且写入不抛异常test_fail_fast_on_non_transient_errors断言非暂态错误被包装为RuntimeError且原因链中保留了原始异常。返回结果结构write_turbopuffer返回一个 DataFrame其中每一行对应一次写入尝试包含write_responses列python类型见 turbopuffer_data_sink.py#L132。每次写入结果包含成功时{status: success, response: Turbopuffer 写入响应}失败时{status: failed, error: 错误信息, rows_not_written: 未写入行数}。内部还通过WriteResult统计了bytes_written与rows_written供执行引擎汇总进度与计量。进阶实战面向生产的多 GPU 向量化写入如果处理的是百万级文本文档推荐参考 docs/examples/text-embeddings.md 中完整的多 GPU 流水线方案其核心链路为用daft.read_parquet从 S3 读取大规模文本数据通过daft.cls类 UDF 用 spaCy 做句子级分块用 Sentence Transformers 在 GPU 上bfloat16精度批量计算向量用col(url).right(50) - col(chunk_id)构造唯一文档 ID调用write_turbopuffer分布式写入关键参数组合如下.write_turbopuffer( namespacetext-embeddings-example, regionaws-us-west-2, id_columnid, vector_columnembedding, distance_metriccosine_distance )该示例展示了id_column/vector_column的真实用途向量列名是embedding而非默认的vectorID 由 URL 与分块编号拼接而成均通过参数显式映射无需预先重命名列。同时该文档强调Daft 会将网络 I/O、CPU 分块与 GPU 推理流水线化并行执行从而在多 GPU 集群上逼近满利用率。小结Daft 对 Turbopuffer 的写入支持可以归纳为三层能力一行一文档的自动映射配合id_column/vector_column重命名机制几乎无需预整形数据参数三通道namespace/region/distance_metric/schema等常用参数直接暴露client_kwargs与write_kwargs透传底层能力兼顾简单与灵活健壮的错误语义暂态错误由客户端重试并优雅降级、非暂态错误快速失败配合结构化写入结果适合嵌入生产调度系统。从 Hugging Face 数据集、S3 中的 Parquet 到本地的任意 DataFrame只需一行.write_turbopuffer(...)即可把数据连同向量送入可检索的搜索引擎。更多自定义选项可以查看 DataFrame.write_turbopuffer 的 API 文档以及 Turbopuffer 连接器的入口文档 docs/connectors/index.md。【免费下载链接】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),仅供参考
返回列表