十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

dlt 集成 LanceDB 目标实战指南:从 Embedding 配置到向量检索与孤儿数据清理

dlt 集成 LanceDB 目标实战指南:从 Embedding 配置到向量检索与孤儿数据清理 dlt 集成 LanceDB 目标实战指南从 Embedding 配置到向量检索与孤儿数据清理【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本指南以 dlt 官方文档 docs/website/docs/dlt-ecosystem/destinations/lancedb.md 为核心骨架结合仓库源码configuration.py、lancedb_adapter.py、jobs.py、utils.py与测试用例tests/load/lancedb深入讲解如何在 dlt 流水线中把数据写入 LanceDB为指定字段自动生成 Embedding 向量、执行向量检索以及使用 merge 写策略时的孤儿块orphan chunks自动清理机制。LanceDB 目标概览LanceDB 中加载数据例如把关系型数据库中的文档表、RAG 语料等同步为 LanceDB 中的向量表。从源码看该目标的能力定义位于 factory.py 的_raw_capabilities()首选加载文件格式为parquet同时支持reference格式支持的最大标识符长度为 200、最大列标识符长度 1024不支持 DDL 事务支持 replace 策略truncate-and-insert支持 merge 策略upsert与insert-only支持嵌套类型用于承载用户自带的向量列并启用use_compliant_nested_typeFalse以保持 Arrow 兼容的嵌套类型存储。安装与模型提供商选择选择 Embedding 模型提供商在使用向量能力前需要先决定使用哪个 embedding 模型提供商。LanceDB 支持的全部提供商以官方文档为准dlt 侧通过TEmbeddingProvider字面量类型configuration.py将其限定为以下集合gemini-textbedrock-textcoheregte-textimagebindinstructoropen-clipopenaisentence-transformershuggingfacecolbertollama安装 dlt 与 lancedb 扩展使用 LanceDB 目标时需要按lancedbextra 安装 dltpip install dlt[lancedb]需要注意lancedbextra 只安装dlt与lancedb本体。模型提供商的 SDK例如 openai、huggingface 客户端等需要根据所选提供商自行安装具体需要哪些库请参照 LanceDB 官方文档中对应提供商的说明。配置目标连接在 dlt 的 secrets 文件默认位于~/.dlt/secrets.toml中添加如下配置段[destination.lancedb] lance_uri .lancedb embedding_model_provider ollama embedding_model mxbai-embed-large embedding_model_provider_host http://localhost:11434 # 可选自定义提供商端点 [destination.lancedb.credentials] api_key api_key # 连接 LanceDB Cloud 的 API Key使用 LanceDB OSS 时注释掉 embedding_model_provider_api_key embedding_model_provider_api_key # 无需鉴权的提供商ollama、sentence-transformers无需填写参数说明与 configuration.py 中的配置规格一致lance_uriLanceDB 实例的位置。未提供时默认为本地磁盘实例。可用 schema 包括/path/to/database本地数据库与db://host:portLanceDB Cloud 远程数据库。api_key连接 LanceDB Cloud 的 API Key使用 LanceDB OSS 时无需提供。embedding_model_provider生成 embedding 的提供商默认值为cohere。embedding_model提供商用于生成 embedding 的具体模型默认值为embed-english-v3.0见 configuration.py。可用选项需查阅对应提供商文档。embedding_model_provider_host支持自定义端点的提供商如 Ollama的完整主机地址需带协议与端口如http://localhost:11434。不指定时使用提供商的默认端点。embedding_model_provider_api_keyembedding 提供商自己的 API Key。对无需鉴权的提供商如 Ollama、sentence-transformers可不填。embedding_model_dimensions源码新增embedding 向量的维度多数情况下由 LanceDB 自动推断仅在少数场景下需要显式指定且必须与所用模型的实际维度一致。源码中LanceDBCredentialsconfiguration.py还包含region默认us-east-1LanceDB Cloud 区域、host_override、client_config与storage_options等字段其中storage_options用于传入 Rust 对象存储凭据。本地数据库命名规则lancedb数据库的命名规则与duckdb相同默认情况下数据库文件名称为pipeline_name.lancedb放置在当前工作目录对于命名目标named destination数据库文件名为destination name.lancedb:pipeline:形式的lance_uri会将数据库文件放在 pipeline 的工作目录中。从源码看on_resolved()configuration.py在配置解析阶段会优先采用显式传入的lance_uri若未提供或提供的是相对本地路径则通过make_location生成name.lancedb形式的本地路径并把最终lance_uri回写进凭据对象。配置 LanceDB Cloudlance_uri以db://schema 开头时被视为 LanceDB Cloud 位置此时必须提供api_key才能连接。dlt 使用与lancedb.connect()函数一致的参数名[destination.lancedb.credentials] api_key api_key region us-east-1 read_consistency_interval2.5其中read_consistency_interval默认是None无读一致性dlt 假定对特定表只有单一写入者。在源码中该值会被转换为timedelta后传给lancedb.connect()configuration.py。提示可以在 credentials 中传入storage_options以便把 LanceDB 数据存储到对象存储桶上。这是 LanceDB 官方支持的能力但 dlt 官方表示尚未对其做过完整测试建议自行验证后再用于生产。定义数据源并写入 LanceDB最小示例加载电影数据并生成向量import dlt from dlt.destinations.adapters import lancedb_adapter movies [ { id: 1, title: Blade Runner, year: 1982, }, { id: 2, title: Ghost in the Shell, year: 1995, }, { id: 3, title: The Matrix, year: 1999, }, ] pipeline dlt.pipeline( pipeline_namemovies, destinationlancedb, ) info pipeline.run( lancedb_adapter( movies, embedtitle, ), table_namemovies, )数据加载后即进入 LanceDB。要使用向量检索必须通过lancedb_adapter包装数据或 dlt 资源明确指定要对哪些字段生成 embedding。上述示例中title列会使用配置的 embedding 提供商与模型生成向量。从实现上看lancedb_adapterlancedb_adapter.py会把被 embed 的列打上x-lancedb-embed表提示VECTORIZE_HINT并确保该列可空nullable因为 Lance 会自行覆盖该值同时写入x-lancedb-remove-orphans表提示。目标端在创建表时schema.py会为这些列追加一个vector向量字段pa.list_(pa.float32(), vec_size)维度取自embedding_model_dimensions或 embedding 函数推断出的ndims()。关于 dataset_name 的说明在上述示例中pipeline 未指定 dataset 名称数据按预期存储在movies表中。如果指定了 dataset 名称dlt 会沿用与其他无 schema 存储相同的模式为所有表加上database_name前缀。例如pipeline dlt.pipeline( pipeline_namemovies, destinationlancedb, dataset_namemovies_db, )此时表名会变为movies_db___movies其中___3 个下划线是可配置的分隔符对应配置项dataset_separator默认___。相关实现见 lancedb_client.py 的make_qualified_table_name。使用适配器指定要向量化的列默认情况下 LanceDB 只作为一个普通数据库使用。要启用其 embedding 能力需要在 dlt 资源中指定要嵌入的字段。lancedb_adapter就是为此提供的辅助函数from dlt.destinations.adapters import lancedb_adapter lancedb_adapter(data, embedtitle)参数说明datadlt 资源对象或 Python 数据结构如字典列表embed需要生成 embedding 的字段名可以是单个字符串或字符串列表。返回值是可直接传给pipeline.run()的 dlt 资源对象。多个字段示例from dlt.destinations.adapters import lancedb_adapter lancedb_adapter( resource, embed[title, description], )源码层面lancedb_adapter.py要求至少提供embed、merge_key、no_remove_orphans之一否则抛出ValueErrorembed必须是字符串或字符串列表否则同样抛出ValueError。必须应用到资源而非整个 source使用lancedb_adapter时要直接应用到资源resource上而不是整个 source。例如从 SQL 数据库加载多个表from dlt.sources.sql_database import sql_database from dlt.destinations.adapters import lancedb_adapter products_tables sql_database().with_resources(products, customers) pipeline dlt.pipeline( pipeline_namepostgres_to_lancedb_pipeline, destinationlancedb, ) # 对需要的资源分别应用适配器 lancedb_adapter(products_tables.products, embeddescription) lancedb_adapter(products_tables.customers, embedbio) info pipeline.run(products_tables)lancedb_adapter内部通过get_resource_for_adapter将原始数据包装成资源再调用resource.apply_hints()注入表提示与列提示lancedb_adapter.py因此对已加载decorated的资源应用也是安全的。使用 Arrow 或 Pandas 加载数据dlt 与 LanceDB 都原生支持 Arrow 与 Pandas因此可以以高性能方式摄取数据无需不必要的重写与拷贝具体做法参考 Arrow/Pandas 数据加载指南。需要注意如果计划使用merge写策略请记得为 Arrow 表启用 load ids 跟踪相关说明见 verified-sources 文档 中“为表添加_dlt_load_id和_dlt_id”一节。这一限制在源码的verify_schema中也有体现lancedb_client.py启用孤儿清理且使用 merge 策略时若表中缺少_dlt_load_id列会抛出DestinationTerminalException并提示通过NORMALIZE__PARQUET_NORMALIZER__ADD_DLT_LOAD_IDTRUE或 config.toml 等价配置开启。访问已加载的数据方式一自建 LanceDB 客户端并注入 pipeline你可以自己创建 LanceDB 客户端将其传给 dlt pipeline 用于加载随后再用同一个客户端查询import dlt import lancedb db lancedb.connect(movies.db) pipeline dlt.pipeline( pipeline_namemovies, destinationdlt.destinations.lancedb(credentialsdb), ) ... tbl db.table(movies) print(tbl.query(magic dog))源码中LanceDBCredentials.parse_native_representationconfiguration.py会识别原生lancedb.DBConnection实例将其存入内部连接槽位并把uri标记为:external:表示由外部客户端持有连接身份对应data_location()中:external:的处理分支configuration.py。方式二从 pipeline 获取已认证客户端也可以通过 pipeline 的目标客户端拿到经过认证的 DB 连接import dlt from lancedb import DBConnection pipeline dlt.pipeline( pipeline_namemovies, destinationlancedb, ) ... with pipeline.destination_client() as job_client: # type: ignore db: DBConnection job_client.db_client # type: ignore tbl db.open_table(movies) tbl.create_scalar_index(id)这里job_client.db_client正是LanceDBClient在初始化时通过credentials.get_conn()建立的连接lancedb_client.py。LanceDBClient还提供了query_table()方法封装 LanceDB 的向量检索lancedb_client.py。自带向量Bring Your Own Vector默认情况下dlt 会根据lancedb_adapter中指定的字段自动添加一个向量列。你也可以选择显式传入向量数据。目前该能力仅在生成具有正确 schema 的 Arrow 表时可用且向量必须声明为固定长度import pyarrow as pa import numpy as np import dlt vector_dim 5 vectors [np.random.rand(vector_dim).tolist() for _ in range(4)] table pa.table( { id: pa.array(list(range(1, 5)), pa.int32()), vector: pa.array( vectors, pa.list_(pa.float32(), vector_dim) ), } ) print(dlt.run(table, table_namevectors, destinationlancedb))在 schema 构建阶段schema.py如果目标表中已存在用户提供的vector列即vector_field_name默认vectordlt 不会再自动追加向量字段而是记录日志要求该 Arrow 列类型必须与向量维度匹配。仓库测试 tests/load/lancedb/test_pipeline.py 的test_bring_your_own_vector覆盖了该路径。写策略Write DispositionLanceDB 目标支持所有 写策略。Replacereplace策略会用资源中的数据替换目标中的既有数据from dlt.destinations.adapters import lancedb_adapter movies [{id: 1, title: Blade Runner, year: 1982}, ...] info pipeline.run( lancedb_adapter( movies, embedtitle, ), write_dispositionreplace, )目标端能力声明支持truncate-and-insert替换策略factory.py加载时会通过tbl.add(records)写入utils.py。Mergemerge写策略基于唯一标识把资源数据与目标数据合并。LanceDB 目标支持upsert与insert-only两种 merge 策略upsert更新已存在的记录并插入新记录insert-only只插入新记录、不更新已有记录见 insert-only 策略说明。可以在资源或适配器中指定 merge 策略、主键与 merge keyfrom typing import Generator from dlt.common.typing import DictStrAny from dlt.destinations.adapters import lancedb_adapter dlt.resource( primary_key[doc_id, chunk_id], merge_key[doc_id], write_disposition{disposition: merge, strategy: upsert}, ) def my_rag_docs( data: list[DictStrAny], ) - Generator[list[DictStrAny], None, None]: yield data pipeline.run( lancedb_adapter( my_rag_docs, merge_keydoc_id ), write_disposition{disposition: merge, strategy: upsert}, primary_key[doc_id, chunk_id], )关键约束文档原话primary_key唯一标识每条记录通常由文档 ID 与块 ID 组成merge_key不能是复合键应对应向量数据库中规范化的doc_id代表数据模型中的文档标识merge_key必须是primary_key的第一个元素该merge_key对 merge 操作期间的文档识别与孤儿清理至关重要保证记录标识正确且与向量数据库概念一致。源码实现上merge 加载由 jobs.py 解析 merge key 与策略并在 utils.py 中分别构造merge_insert(merge_key).when_matched_update_all().when_not_matched_insert_all()upsert或when_not_matched_insert_all()insert-only操作。孤儿清理Orphan Removal在 merge 操作中当更新或删除父文档时LanceDB 会自动清理孤儿块orphaned chunks。若要关闭该特性from dlt.destinations.adapters import lancedb_adapter movies [{id: 1, title: Blade Runner, year: 1982}, ...] pipeline.run( lancedb_adapter( movies, embedtitle, no_remove_orphansTrue # 通过 no_remove_orphans 标志关闭 ), write_disposition{disposition: merge, strategy: upsert}, primary_key[doc_id, chunk_id], )虽然可以省略merge_key此时假定其为primary_key的第一个元素但为清晰起见建议同时显式指定两者。实现细节no_remove_orphansTrue会通过适配器写入x-lancedb-remove-orphans表提示lancedb_adapter.py。若未关闭且 merge 策略非 insert-only客户端在表链完成加载后会追加一个LanceDBRemoveOrphansJoblancedb_client.py对于根表按规范化doc_id来自 merge key 或 primary_key 首元素utils.py删除与本次 load id 不一致的旧记录对于嵌套表则依据_dlt_root_id删除 payload 中不存在的孤儿行jobs.py最终通过when_not_matched_by_source_delete(delete_condition)执行删除utils.py。注意孤儿清理依赖_dlt_id与_dlt_load_id字段而 Arrow 表加载时默认并不包含这两个字段。需要在 normalize 配置中将add_dlt_id选项设为true以启用见 Arrow/Pandas 指南。verify_schema也会在缺少_dlt_load_id时直接报错提示lancedb_client.py。Append这是默认策略将数据追加到目标中的既有数据之后。加载时对应tbl.add(records)utils.py。其他目标选项dataset_separator用于分隔 dataset 名称与表名的字符默认___。vector_field_name存储向量 embedding 的特殊字段名默认vector。max_retriesembedding 操作的最大重试次数设为 0 可禁用重试默认 3。对应 configuration.py 中LanceDBClientOptions.max_retries的说明LanceDB 的EmbeddingFunction在失败后会以指数退避方式重试请求dlt 将其透传给 embedding 函数lancedb_client.py。sentinel_table_name源码新增默认dltSentinelTable作为封装 dataset 的哨兵表因为 LanceDB 没有 schema 概念该表充当代理将相关 dlt 表分组configuration.py。always_refresh_views源码新增每次查询前重建视图默认False。环境变量注入部分 embedding 提供商通过环境变量读取 API Key。dlt 在创建 embedding 函数前会调用set_non_standard_providers_environment_variablesutils.py把embedding_model_provider_api_key注入对应环境变量提供商环境变量cohereCOHERE_API_KEYgemini-textGOOGLE_API_KEYopenaiOPENAI_API_KEYhuggingfaceHUGGINGFACE_API_KEY因为 LanceDB 没有标准化的跨提供商 API Key 注入方式有些用环境变量有些接受参数dlt 统一设置环境变量来兼容lancedb_client.py。限制与注意事项dbt 支持LanceDB 目标不支持 dbt 集成。dlt 状态同步LanceDB 目标支持dlt状态的同步通过WithStateSync与get_stored_state实现lancedb_client.py。merge 约束启用孤儿清理时不允许复合 merge keyverify_schema中会抛出DestinationTerminalExceptionlancedb_client.py。数据类型一致性merge 时 LanceDB 要求源数据 schema 与目标表 schema 完全一致列名、顺序、类型dlt 在写入前会为目标 schema 补齐向量列utils.py。存储桶选项storage_options可用于把数据放到对象存储但 dlt 官方尚未完整测试。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表