Parlant Qdrant 向量数据库适配器:用持久化向量存储替换默认内存存储的完整实践
【免费下载链接】parlantBuild reliable customer-facing AI agents with Parlant: an interaction control harness optimized for controlled, consistent, and predictable LLM interactions.项目地址: https://gitcode.com/GitHub_Trending/pa/parlant
本文基于 Parlant 仓库中的 Qdrant 适配器文档 与 Qdrant 适配器源码 展开,讲解如何在 Parlant 中通过QdrantDatabase将术语表、固定回复、能力、旅程等四类向量存储从默认内存(transient)存储迁移到生产级持久化存储。读完本文,你将掌握 Qdrant 适配器的完整接入代码、集合命名与自动同步机制、Where 过滤器到 Qdrant Filter 的映射规则,以及 Windows 文件锁、嵌入器切换、数据持久化验证等常见问题的排查方法。
为什么需要 Qdrant 适配器
Parlant 默认的向量存储由 TransientVectorDatabase 提供,底层基于nano_vectordb在内存中维护向量集合:
# src/parlant/adapters/vector_db/transient.py(节选) self._databases[name] = nano_vectordb.NanoVectorDB(embedder.dimensions)这种实现适合开发调试,但服务重启后向量数据即丢失。Qdrant 适配器通过QdrantDatabase类实现了核心抽象 VectorDatabase 接口(create_collection、get_collection、get_or_create_collection、delete_collection、upsert_metadata、remove_metadata、read_metadata),将向量数据的持久化能力无缝接入 Parlant 的容器(Container)体系,同时保持与现有向量存储接口的完全兼容——上层只需在configure_container中替换 Store 实现,无需修改任何业务逻辑。
前置条件
安装 Qdrant 适配器(以可选依赖方式安装):
pip install parlant[qdrant]从 pyproject.toml 可见,
qdrantextra 实际拉取的是qdrant-client>=1.7.0。选择存储方式:本地文件系统,或 Qdrant Cloud 远程集群。
Python 版本:文档中标注为 Python 3.8+,但从当前仓库的 pyproject.toml 看,Parlant 实际要求
requires-python = ">=3.10,<3.15"(该限制是为兼容 torch 2.8+ 与 triton),因此请以 Python 3.10+ 为实际前提。本地存储模式需要一个可写目录;Cloud 模式则需要一个 Qdrant Cloud 账户(URL 与 API Key)。
快速上手:在 configure_container 中接入 Qdrant
核心思路是:在 SDK 的configure_container钩子中创建QdrantDatabase实例,然后把四个向量 Store(GlossaryStore、CannedResponseStore、CapabilityStore、JourneyStore)全部替换为对应 VectorStore 实现。以下完整代码继承自 官方文档:
import parlant.sdk as p from pathlib import Path from contextlib import AsyncExitStack from parlant.adapters.vector_db.qdrant import QdrantDatabase from parlant.core.nlp.embedding import EmbedderFactory, EmbeddingCache, Embedder from parlant.core.loggers import Logger from parlant.core.nlp.service import NLPService from parlant.core.glossary import GlossaryVectorStore, GlossaryStore from parlant.core.canned_responses import CannedResponseVectorStore, CannedResponseStore from parlant.core.capabilities import CapabilityVectorStore, CapabilityStore from parlant.core.journeys import JourneyVectorStore, JourneyStore from parlant.adapters.db.transient import TransientDocumentDatabase async def configure_container(container: p.Container) -> p.Container: embedder_factory = EmbedderFactory(container) async def get_embedder_type() -> type[Embedder]: return type(await container[NLPService].get_embedder()) exit_stack = AsyncExitStack() qdrant_db = await exit_stack.enter_async_context( QdrantDatabase( logger=container[Logger], path=Path("./qdrant_data"), embedder_factory=EmbedderFactory(container), embedding_cache_provider=lambda: container[EmbeddingCache], ) ) # For Qdrant Cloud, replace the above with: # qdrant_db = await exit_stack.enter_async_context( # QdrantDatabase( # logger=container[Logger], # url="https://your-cluster-id.us-east4-0.gcp.cloud.qdrant.io", # api_key="your-api-key-here", # embedder_factory=EmbedderFactory(container), # embedding_cache_provider=lambda: container[EmbeddingCache], # ) # ) # Configure stores using vector database container[GlossaryStore] = await exit_stack.enter_async_context( GlossaryVectorStore( id_generator=container[p.IdGenerator], vector_db=qdrant_db, document_db=TransientDocumentDatabase(), embedder_factory=embedder_factory, embedder_type_provider=get_embedder_type, ) # type: ignore ) container[CannedResponseStore] = await exit_stack.enter_async_context( CannedResponseVectorStore( id_generator=container[p.IdGenerator], vector_db=qdrant_db, document_db=TransientDocumentDatabase(), embedder_factory=embedder_factory, embedder_type_provider=get_embedder_type, ) # type: ignore ) container[CapabilityStore] = await exit_stack.enter_async_context( CapabilityVectorStore( id_generator=container[p.IdGenerator], vector_db=qdrant_db, document_db=TransientDocumentDatabase(), embedder_factory=embedder_factory, embedder_type_provider=get_embedder_type, ) # type: ignore ) container[JourneyStore] = await exit_stack.enter_async_context( JourneyVectorStore( id_generator=container[p.IdGenerator], vector_db=qdrant_db, document_db=TransientDocumentDatabase(), embedder_factory=embedder_factory, embedder_type_provider=get_embedder_type, ) # type: ignore ) return container async def main(): async with p.Server(configure_container=configure_container) as server: agent = await server.create_agent( name="My Agent", description="Agent using Qdrant for persistent storage", ) # Test: Create a term to verify Qdrant is working term = await agent.create_term( name="Example Term", description="This is stored in Qdrant", ) print(f"Created term: {term.name}") # All vector operations now use Qdrant关键参数说明
结合 QdrantDatabase 构造函数,各参数含义如下:
| 参数 | 说明 |
|---|---|
logger | 从容器中取Logger,用于记录超时重试等警告。 |
tracer | 源码中还有tracer: Tracer参数,用于链路追踪;文档示例未显式传入时以容器默认行为为准。 |
path | 本地 Qdrant 的数据目录(如Path("./qdrant_data")),启动后会在该目录下生成 Qdrant 数据库文件。与url互斥:传path走本地模式。 |
url/api_key | Qdrant Cloud 集群地址与密钥。走远程模式时,客户端会显式设置60 秒超时(见下文源码分析),以容忍大批量操作和较慢的网络。 |
embedder_factory | EmbedderFactory,用于按嵌入器类型创建 embedder 并决定向量维度。 |
embedding_cache_provider | 嵌入缓存提供者(通常注入container[EmbeddingCache]),避免对相同文本重复调用嵌入模型。 |
| 都不传 | 从源码看,path和url均未提供时,QdrantClient会以":memory:"初始化,仅适合测试场景(源码)。 |
四个 VectorStore 的构造参数完全一致,其中:
id_generator=container[p.IdGenerator]:统一使用容器中的 ID 生成器;vector_db=qdrant_db:即上面的 Qdrant 数据库实例;document_db=TransientDocumentDatabase():注意这里文档(非向量)存储仍可使用内存实现,持久化的只是向量部分;如需完整持久化,可另行接入文档数据库适配器;embedder_type_provider=get_embedder_type:一个异步回调,返回当前 NLP 服务的嵌入器类型,用于决定集合命名与维度。
生命周期管理要点:QdrantDatabase与四个 VectorStore 都是异步上下文管理器,必须通过AsyncExitStack的enter_async_context注册,确保服务关闭时按序退出、释放文件句柄(这一点在 Windows 上尤其重要,见下文)。
源码深潜:Qdrant 适配器内部机制
1. 双集合结构与命名规则
每个逻辑集合在 Qdrant 中实际对应两个物理集合(create_collection 源码):
- 嵌入集合
{name}_{EmbedderTypeName}:存放真实向量,尺寸取embedder.dimensions,距离度量固定为 COSINE; - 未嵌入集合
{name}_unembedded:向量尺寸仅为 1(占位向量[0]),作为文档的SSOT(单一事实来源),存储完整 payload 与校验和(checksum)。
例如嵌入器为OpenAITextEmbedding3Large时,你会看到glossary_OpenAITextEmbedding3Large、glossary_unembedded这样的集合名——这正是文档“Check Collections”一节中列出集合名的来源。两个集合都会在id字段上创建 KEYWORD 类型的payload index(_ensure_payload_index)。
文档 ID 是字符串,而 Qdrant point ID 支持整数,适配器用 SHA-256 哈希将其映射到 int64 安全范围内(_string_id_to_int):
hash_value = int(hashlib.sha256(doc_id.encode()).hexdigest()[:15], 16) return hash_value % (2**63 - 1)2. 写入与自动同步:unembedded 到 embedded 的迁移
get_collection/get_or_create_collection打开集合时会调用_load_collection_documents:先用document_loader把 unembedded 集合里的旧文档逐一加载、按需迁移为新 schema,再触发_index_collection完成向量重建:
- 按 payload 中的文档
id建立两个集合的映射; - 删除 unembedded 中已不存在的旧向量点;
- 只有当
checksum变化时才重新调用嵌入模型,否则直接复用旧向量,避免不必要的模型调用; - 加载失败的文档会被写入独立的
failed_migrations集合以便调试(源码); - 通过
metadata集合中的版本键({collection_name}_version)判断是否需要重建索引。
这就是文档中“Collection Sync:集合在嵌入器或 schema 变化时自动同步;大集合首次访问可能较慢”的底层实现。
3. 超时重试与远程超时
写入操作(insert_one/update_one)通过_retry_on_timeout_async包裹,针对超时类错误做**最多 3 次、指数退避(1s、2s、4s)**的重试;同时所有阻塞式 QdrantClient 调用都通过asyncio.to_thread放到线程池执行,避免阻塞事件循环。远程模式下客户端初始化时设置 60 秒超时(源码)。
4. Where 过滤器到 Qdrant Filter 的映射
上层 Store 使用类 MongoDB 的Where字典过滤器,适配器通过_convert_where_to_qdrant_filter将其翻译为 Qdrant 原生过滤条件:
| Parlant Where 操作符 | Qdrant Filter 形式 |
|---|---|
$eq | FieldCondition(match=MatchValue(...))放入must |
$ne | must_not+MatchValue |
$gt/$gte/$lt/$lte | FieldCondition(range=Range(...))放入must |
$in | FieldCondition(match=MatchAny(...))放入must |
$nin | must_not+MatchAny |
$and | 递归转换为多个子条件放入must |
$or | 递归转换为多个子条件放入should |
另外,执行find/find_one/update_one/delete_one/ 相似度检索前,适配器会递归提取过滤器中出现的所有字段名并确保对应 payload index 存在(源码);若服务端过滤失败,则回退为全量 scroll + 内存过滤(matches_filters)保证正确性(源码)。
5. 相似度检索与距离换算
do_find_similar_documents先用嵌入器将查询文本向量化(空向量直接返回空结果并打警告),再调用query_points执行带过滤的 ANN 检索,最后把 Qdrant 的余弦相似度得分换算为距离:
SimilarDocumentResult(document=..., distance=1.0 - result.score)读写并发由ReaderWriterLock控制:查询取 reader 锁,写操作取 writer 锁。元数据(如各 Store 的 schema 版本VectorDocumentStoreMigrationHelper.get_store_version_key(...))则统一存放在名为metadata的集合中,以单个固定 point(__metadata__)承载整个键值文档(源码)。
6. Windows 文件锁处理
本地模式下,__aenter__在 Windows(sys.platform == "win32")上会把打开重试次数从 1 提升到 5,遇到already accessed错误时按 0.05s 递增退避重试(源码);__aexit__则显式close()客户端、释放集合引用,并在 Windows 上触发gc.collect()加 50ms 等待,确保操作系统及时释放文件锁(源码)。这就是文档中“Windows File Locks:使用async with上下文管理器,适配器会自动处理文件锁重试”的具体实现。
验证 Qdrant 集成是否生效
检查集合
Qdrant Cloud:集合会出现在你的 Qdrant 控制台,典型名称包括:
glossary_OpenAITextEmbedding3Largeglossary_unembeddedcapabilities_OpenAITextEmbedding3Largecanned_responses_OpenAITextEmbedding3Large
本地 Qdrant:指定路径下会生成包含 Qdrant 数据库文件的文件夹。
确认没有回退到内存存储
Qdrant 配置正确时:
parlant-data文件夹中不会出现向量文件(Parlant 服务端的默认数据目录即parlant-data,见 server.py 中DEFAULT_HOME_DIR的定义);- 向量数据仅存储在 Qdrant(本地或云端);
- 数据在服务重启后依然存在。
测试向量检索与持久化
term = await agent.create_term( name="Test Term", description="This should be stored in Qdrant", ) # Then chat with agent about "test term" - it should understand via vector search # Test persistence: close the server and run again # The term should still be available after restart仓库中的 tests/adapters/vector_db/test_qdrant.py 对同一机制做了完整覆盖,可作为验证清单参考:test_loading_collections验证关闭数据库后重新打开仍能取回文档(持久化);test_find_similar_documents验证 ANN 检索命中预期文档;test_and_operator_with_multiple_conditions、test_or_operator_with_multiple_conditions、test_that_in_filter_works_with_list_of_strings分别验证$and/$or/$in过滤器;test_that_documents_are_indexed_when_changing_embedder_type验证切换嵌入器后集合自动重建索引。
常见问题与故障排查
集成未生效(仍在使用内存存储)
症状:Qdrant 控制台中看不到集合;parlant-data文件夹里出现了向量数据;重启服务后数据丢失。
解决:确认configure_container中四个向量 Store 均已替换为 VectorStore 实现并注入qdrant_db;确认使用AsyncExitStack正确管理QdrantDatabase与各 Store 的生命周期。
Windows 文件锁
在 Windows 上使用async with上下文管理器,适配器会自动处理文件锁重试(见上文源码分析)。
集合同步
集合会在嵌入器或 schema 变化时自动同步;大集合首次访问(触发重建索引)可能需要较长时间,这是预期行为。
更换嵌入器
切换嵌入器类型后,旧嵌入器命名的集合(如glossary_OldEmbedder)会一直保留,直到手动删除——因为集合名中包含嵌入器类型名,新集合会按新名字创建,旧数据不会自动清理。
性能
生产环境建议使用 Qdrant Cloud 或 Qdrant 服务端:本地模式不支持 payload 索引。使用本地 Qdrant 时你会看到相关警告,这是预期行为,可以忽略。其他建议:
- 使用嵌入缓存(
embedding_cache_provider)减少对嵌入模型的重复调用; - 使用 Qdrant Cloud/服务端以获得 payload 索引支持;
- 考虑拆分过大的集合。
连接问题
- 本地:确认路径存在且可写;
- 远程:核对 URL 与 API Key。
数据未持久化
- 检查文件路径是否正确、可写;
- 远程部署时核对连接配置;
- 通过“关闭服务再重启,数据仍在”这一行为做最终验证。
小结
| 能力 | 说明 |
|---|---|
| 持久化存储 | 以 Qdrant 替换基于nano_vectordb的默认内存存储,重启不丢数据 |
| 自动同步 | 嵌入器或 schema 变化时自动重建嵌入集合,checksum 未变的文档复用旧向量 |
| 迁移容错 | 加载失败的文档落入failed_migrations集合,便于排查 |
| Windows 支持 | 文件锁自动重试、显式关闭与 GC 处理 |
| 过滤器支持 | $eq/$ne/$gt/$gte/$lt/$lte/$in/$nin/$and/$or全量映射,含内存过滤回退 |
| 依赖与前提 | pip install parlant[qdrant](qdrant-client>=1.7.0),Python 3.10+(仓库requires-python = ">=3.10,<3.15"),本地可写目录或 Qdrant Cloud 账户 |
实现与验证材料索引:适配器实现 src/parlant/adapters/vector_db/qdrant.py、核心抽象 src/parlant/core/persistence/vector_database.py、默认内存实现 src/parlant/adapters/vector_db/transient.py、依赖声明 pyproject.toml、测试 tests/adapters/vector_db/test_qdrant.py、原始文档 docs/adapters/vector_db/qdrant.md。
【免费下载链接】parlantBuild reliable customer-facing AI agents with Parlant: an interaction control harness optimized for controlled, consistent, and predictable LLM interactions.项目地址: https://gitcode.com/GitHub_Trending/pa/parlant
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考