LangChain与Milvus向量数据库DML实战:数据更新与删除操作指南

发布时间:2026/8/25 21:15:08
LangChain与Milvus向量数据库DML实战:数据更新与删除操作指南 在构建基于大语言模型的应用时向量数据库是连接非结构化文本与模型理解能力的关键桥梁。然而很多开发者在完成向量存储后往往忽略了数据管理的重要性导致应用难以维护和迭代。本文将聚焦于 LangChain 与 Milvus 的深度集成手把手带你完成从数据插入、更新到删除的完整 DML数据操作语言实战确保你的 AI 应用数据层既高效又健壮。无论你是刚接触向量数据库的新手还是希望优化现有 LangChain 工作流的开发者本文都将提供一套可直接复用的代码方案和清晰的工程实践思路。1. 背景与核心概念为什么需要关注 DML在 LangChain 生态中Milvus 因其高性能、可扩展性以及对海量向量数据的友好支持成为了许多开发者的首选向量数据库。我们通常关注如何将文档切分、嵌入并存入 Milvus即“写”操作但这仅仅是开始。一个成熟的 AI 应用其知识库必然是动态的数据更新源文档内容修订后需要同步更新向量库中的对应条目。数据删除清理过时、无效或敏感的数据。条件操作需要根据元数据如文档ID、类别、创建时间进行批量更新或删除。这些操作统称为 DML。在 Milvus 中DML 主要围绕insert、upsert、delete等操作展开。LangChain 的Milvus向量存储封装类虽然提供了基础的add_documents和delete方法但要实现灵活、精准的数据管理我们需要深入其底层连接直接使用 Milvus SDK 的强大功能。本文将深入探讨如何结合 LangChain 的高级抽象与 Milvus SDK 的精细控制构建一个可维护的向量数据管理流程。2. 环境准备与版本说明在开始实战之前请确保你的开发环境已就绪。以下版本为本文撰写时的稳定版本实际操作时请根据你的项目需求进行微调。核心组件版本Python: 3.8LangChain: 0.1.x (本文示例基于0.1.10的语法请注意 LangChain 版本迭代较快部分导入路径可能变化)Milvus: 2.3.x 或以上 (Standalone 或 Cluster 模式均可)pymilvus: 2.3.x (Milvus 的 Python SDK)Embedding 模型: 本文使用text-embedding-ada-002的 OpenAI 接口你也可以替换为sentence-transformers等本地模型。项目依赖 (requirements.txt):langchain0.1.10 langchain-openai0.0.5 # 用于调用OpenAI Embedding pymilvus2.3.6 openai1.12.0 python-dotenv1.0.0 # 用于管理环境变量关键环境变量 (.env文件):# 如果你的Milvus需要认证 MILVUS_URIhttp://localhost:19530 MILVUS_USERusername MILVUS_PASSWORDpassword MILVUS_DB_NAMEdefault # OpenAI (或其他Embedding服务) OPENAI_API_KEYyour_openai_api_key_here启动 Milvus 服务:如果你使用 Docker 运行 Standalone 模式的 Milvus可以使用以下命令docker run -d --name milvus-standalone \ -p 19530:19530 \ -p 9091:9091 \ -v /path/to/milvus/data:/var/lib/milvus \ -v /path/to/milvus/conf:/etc/milvus \ milvusdb/milvus:v2.3.6-standalone启动后确保可以通过19530端口连接到 Milvus。3. 核心原理与 LangChain 集成拆解3.1 LangChain 中 Milvus 向量存储的运作方式LangChain 的Milvus类通常从langchain.vectorstores导入是一个高级封装。它的核心工作流程是接收文档接收Document对象列表每个Document包含page_content和metadata。生成嵌入通过指定的Embeddings模型如OpenAIEmbeddings将page_content转换为向量。与 Milvus 交互通过pymilvusSDK将向量和元数据插入到指定的 Milvus 集合Collection中。创建索引可选在插入数据后为向量字段创建索引以加速检索。Milvus.from_documents()方法封装了创建集合、插入数据、建立索引的完整过程。3.2 DML 操作的底层支撑pymilvus SDK为了实现更精细的 DML 控制我们必须理解其底层依赖——pymilvus。几个关键对象Collection: 代表 Milvus 中的一个数据集合是操作的入口。MutationResult: 执行插入、删除、更新操作后返回的结果包含影响的行数等信息。主键 (Primary Key): Milvus 集合必须有一个主键字段通常是int64或varchar。在 LangChain 默认配置中它会自动生成一个pk字段。为了支持更新和删除我们必须自定义一个有意义且唯一的主键如doc_id并将其存入metadata。3.3 自定义主键的策略这是实现可靠 DML 的基石。LangChain 默认使用自动生成的 ID这不利于追踪。我们的策略是在创建集合时显式定义一个主键字段例如doc_id(VARCHAR类型)。在文档的metadata中必须包含这个doc_id。插入数据时将doc_id作为主键值传入。后续的更新和删除操作都可以通过这个doc_id来精确定位数据。4. 完整实战构建支持 CRUD 的向量数据管理模块接下来我们将一步步构建一个完整的示例。假设我们管理一个“技术文章”知识库每篇文章有唯一ID、标题、内容和分类。4.1 初始化创建带自定义主键的 Milvus 集合首先我们不再完全依赖 LangChain 的自动创建而是先通过pymilvus明确定义集合结构。# file: init_milvus_collection.py from pymilvus import connections, FieldSchema, CollectionSchema, DataType, Collection, utility from langchain_openai import OpenAIEmbeddings import os from dotenv import load_dotenv load_dotenv() # 1. 连接 Milvus connections.connect( aliasdefault, urios.getenv(MILVUS_URI, http://localhost:19530), useros.getenv(MILVUS_USER, ), # 如果未设置认证则为空字符串 passwordos.getenv(MILVUS_PASSWORD, ), db_nameos.getenv(MILVUS_DB_NAME, default) ) # 2. 定义集合名称和 Embedding 维度 collection_name tech_articles_with_pk embedding_dim 1536 # OpenAI text-embedding-ada-002 的维度 # 3. 删除已存在的同名集合仅用于演示生产环境慎用 if utility.has_collection(collection_name): utility.drop_collection(collection_name) # 4. 定义字段 Schema fields [ FieldSchema(namepk, dtypeDataType.VARCHAR, is_primaryTrue, auto_idFalse, max_length100), # 自定义主键不自增 FieldSchema(namedoc_id, dtypeDataType.VARCHAR, max_length100), # 与pk保持一致方便查询 FieldSchema(nametitle, dtypeDataType.VARCHAR, max_length500), FieldSchema(namecontent, dtypeDataType.VARCHAR, max_length65535), FieldSchema(namecategory, dtypeDataType.VARCHAR, max_length50), FieldSchema(nameembedding, dtypeDataType.FLOAT_VECTOR, dimembedding_dim), FieldSchema(nameupdate_time, dtypeDataType.INT64), # 用于记录更新时间戳 ] # 5. 创建集合 Schema schema CollectionSchema(fields, description技术文章向量库支持通过doc_id更新) # 6. 创建集合 collection Collection(namecollection_name, schemaschema) print(f集合 {collection_name} 创建成功主键字段为 pk。) # 7. 为向量字段创建索引HNSW是常用索引 index_params { metric_type: L2, index_type: HNSW, params: {M: 8, efConstruction: 200}, } collection.create_index(embedding, index_params) print(向量索引创建成功。) # 8. 加载集合到内存执行搜索前必须加载 collection.load() print(集合已加载。)4.2 封装数据操作类我们将创建一个工具类封装 LangChain 的文档添加以及基于pymilvus的更新、删除操作。# file: milvus_dml_manager.py from typing import List, Optional, Dict, Any from langchain.schema import Document from langchain.vectorstores import Milvus from langchain_openai import OpenAIEmbeddings from pymilvus import Collection, connections import openai import time import os class MilvusDMLManager: def __init__(self, collection_name: str, embedding_model: OpenAIEmbeddings): 初始化管理器。 :param collection_name: Milvus 集合名称 :param embedding_model: LangChain Embeddings 模型实例 self.collection_name collection_name self.embedding_model embedding_model # 获取 Milvus 集合对象 self.collection Collection(collection_name) # 初始化 LangChain 的 Milvus 向量存储用于检索和基础添加 # 注意这里我们传入已存在的集合名并禁用自动创建schema self.vector_store Milvus( embedding_functionself.embedding_model, collection_nameself.collection_name, connection_args{ uri: os.getenv(MILVUS_URI, http://localhost:19530), user: os.getenv(MILVUS_USER, ), password: os.getenv(MILVUS_PASSWORD, ), db_name: os.getenv(MILVUS_DB_NAME, default) }, auto_schemaFalse, # 关键告诉 LangChain 不要自动创建schema ) def add_documents(self, documents: List[Document]) - List[str]: 使用 LangChain 添加文档并确保主键被正确设置。 LangChain 的 add_documents 会调用 embedding 模型并插入数据。 但我们需要确保 Document 的 metadata 中包含我们定义的 doc_id 并且这个 doc_id 会被用作 Milvus 的主键 pk。 # 在插入前可以验证 documents 是否包含 doc_id for doc in documents: if doc_id not in doc.metadata: raise ValueError(f文档缺失 doc_id metadata: {doc.page_content[:100]}...) # 调用 LangChain 的添加方法。 # Milvus 类内部会处理 embedding 生成并将 metadata 中的所有字段插入。 # 根据我们的集合定义它会把 metadata 中的 doc_id 值同时赋给 doc_id 字段和主键 pk 字段。 # 这是通过在创建集合时将 pk 字段的 auto_id 设为 False 实现的插入时必须提供主键值。 # LangChain 的 Milvus 模块会将 id 或指定的字段作为主键。 # 我们需要确保初始化 Milvus 向量存储时通过 primary_field 参数指定主键字段为 pk。 # 但更直接的方式是在初始化 self.vector_store 时确保其配置与我们的集合匹配。 # 一个更稳妥的做法是直接使用 pymilvus 插入但为了利用 LangChain 的 embedding 流程我们稍作调整 # 实际上LangChain 的 Milvus 类在插入时会尝试将每个 Document 的 id (如果不存在则生成) 作为主键。 # 我们需要将我们的 doc_id 赋值给 Document 的 id 属性。 for doc in documents: doc.id doc.metadata[doc_id] # 添加文档 pks self.vector_store.add_documents(documents) print(f成功添加 {len(pks)} 个文档主键列表: {pks[:5]}...) # 打印前5个 return pks def upsert_document(self, document: Document): 更新或插入文档。 逻辑如果存在相同 doc_id 的文档则先删除再插入Milvus 的 Upsert 在 2.3 版本可用这里提供通用方法。 :param document: LangChain Document 对象其 metadata 必须包含 doc_id doc_id document.metadata.get(doc_id) if not doc_id: raise ValueError(文档必须包含 doc_id metadata 以支持 upsert。) # 1. 生成嵌入向量 embedding self.embedding_model.embed_query(document.page_content) # 2. 准备数据行字段顺序需与集合定义一致 data [ [doc_id], # pk [doc_id], # doc_id [document.metadata.get(title, )], # title [document.page_content], # content [document.metadata.get(category, )], # category [embedding], # embedding [int(time.time())] # update_time ] # 3. 构建删除表达式 expr fpk {doc_id} # 4. 执行删除如果存在 delete_result self.collection.delete(expr) print(f尝试删除 doc_id{doc_id}, 删除条数: {delete_result.delete_count}) # 5. 插入新数据 insert_result self.collection.insert(data) print(f插入 doc_id{doc_id} 成功主键: {insert_result.primary_keys}) # 6. 刷新集合使更改立即可见对于小规模操作可以最后统一刷新 self.collection.flush() return insert_result def delete_by_doc_id(self, doc_id: str) - int: 根据文档ID删除数据。 :param doc_id: 要删除的文档ID :return: 被删除的行数 expr fpk {doc_id} # 或 fdoc_id {doc_id}两者值相同 delete_result self.collection.delete(expr) print(f删除表达式: {expr}, 影响行数: {delete_result.delete_count}) self.collection.flush() return delete_result.delete_count def delete_by_condition(self, condition_expr: str) - int: 根据条件表达式批量删除。 :param condition_expr: Milvus 布尔表达式例如 category deprecated :return: 被删除的行数 # 注意删除前最好先查询确认避免误删 query_result self.collection.query(exprcondition_expr, output_fields[pk, title]) if query_result: print(f即将删除以下 {len(query_result)} 条记录:) for item in query_result[:5]: # 预览前5条 print(f - PK: {item[pk]}, Title: {item.get(title)}) if len(query_result) 5: print(f ... 以及另外 {len(query_result)-5} 条记录) # 这里可以添加人工确认逻辑生产环境建议有审批流程 delete_result self.collection.delete(condition_expr) print(f删除完成影响行数: {delete_result.delete_count}) self.collection.flush() return delete_result.delete_count def search_similar(self, query: str, k: int 3) - List[Dict]: 使用 LangChain 的向量存储进行相似性搜索。 docs self.vector_store.similarity_search(query, kk) results [] for doc in docs: results.append({ content: doc.page_content, metadata: doc.metadata, score: doc.metadata.get(score, 0) # LangChain 可能不直接返回分数 }) return results4.3 主程序执行完整的 DML 操作链现在我们编写一个主程序来使用这个管理器。# file: main_demo.py from milvus_dml_manager import MilvusDMLManager from langchain.schema import Document from langchain_openai import OpenAIEmbeddings import os from dotenv import load_dotenv load_dotenv() # 初始化 Embedding 模型 embeddings OpenAIEmbeddings( modeltext-embedding-ada-002, openai_api_keyos.getenv(OPENAI_API_KEY) ) # 初始化管理器 manager MilvusDMLManager( collection_nametech_articles_with_pk, embedding_modelembeddings ) print( 1. 添加初始文档 ) initial_docs [ Document( page_contentLangChain是一个用于开发由语言模型驱动的应用程序的框架。, metadata{doc_id: doc_001, title: LangChain简介, category: framework} ), Document( page_contentMilvus是一个开源向量数据库专为海量向量相似性搜索而设计。, metadata{doc_id: doc_002, title: Milvus概述, category: database} ), Document( page_contentDML包括INSERT、UPDATE、DELETE等数据操作命令。, metadata{doc_id: doc_003, title: 数据库DML, category: database} ), ] added_ids manager.add_documents(initial_docs) print(f添加的文档ID: {added_ids}\n) print( 2. 相似性搜索测试 ) search_results manager.search_similar(什么是向量数据库, k2) for res in search_results: print(f- 内容: {res[content][:80]}...) print(f 元数据: {res[metadata]}\n) print( 3. 更新文档 (Upsert) ) # 假设 doc_002 的内容需要更新 updated_doc Document( page_contentMilvus是一个云原生开源向量数据库支持秒级检索万亿级向量广泛应用于AI、推荐、搜索等领域。, metadata{doc_id: doc_002, title: Milvus详解, category: vector-db} ) upsert_result manager.upsert_document(updated_doc) print(fUpsert 操作完成。\n) print( 4. 验证更新结果 ) # 再次搜索查看更新后的内容 search_results_after_update manager.search_similar(云原生向量数据库, k1) for res in search_results_after_update: print(f更新后搜索结果: {res[content]}) print(f对应元数据: {res[metadata]}\n) print( 5. 条件删除 ) # 删除分类为 framework 的文档这里会删除 doc_001 deleted_count manager.delete_by_condition(category framework) print(f条件删除完成共删除 {deleted_count} 条记录。\n) print( 6. 按 ID 删除 ) # 删除 doc_003 deleted_count_id manager.delete_by_doc_id(doc_003) print(f按ID删除完成共删除 {deleted_count_id} 条记录。\n) print( 7. 最终搜索验证 ) # 现在集合中应该只剩下 doc_002 final_results manager.search_similar(数据库, k5) print(f剩余文档数量: {len(final_results)}) for res in final_results: print(f- 剩余文档: {res[metadata].get(title)} (ID: {res[metadata].get(doc_id)}))4.4 运行与验证确保 Milvus 服务正在运行。按顺序执行脚本python init_milvus_collection.py python main_demo.py观察控制台输出你应该能看到完整的“添加 - 搜索 - 更新 - 再搜索 - 条件删除 - 按ID删除 - 最终验证”流程。预期输出关键点初始添加成功并返回主键列表。第一次搜索能查到“Milvus是一个开源向量数据库...”。Upsert 后再次搜索“云原生向量数据库”返回的内容应变为更新后的长文本。条件删除后分类为framework的文档被移除。按ID删除后doc_003被移除。最终集合中仅剩更新后的doc_002。5. 常见问题与排查思路在集成 LangChain 与 Milvus 进行 DML 操作时你可能会遇到以下典型问题问题现象可能原因排查思路与解决方案插入失败提示主键冲突1. 主键字段auto_idTrue但插入时又提供了值。2. 多次插入相同的doc_id。1. 检查集合 Schema 定义确保is_primaryTrue的字段其auto_idFalse。2. 实现 Upsert 逻辑先删后插或检查业务逻辑避免重复插入。更新后搜索不到最新数据1. 数据未刷新 (flush)。2. 索引未重建对于大量更新后。3. 搜索时未加载集合。1. 在插入/删除操作后调用collection.flush()。2. 对于大量数据变更考虑在后台异步重建索引。3. 确保在执行搜索前collection.load()。LangChain 的add_documents不生效1.auto_schemaTrue但与已存在集合的 Schema 冲突。2. Document 的metadata字段与集合 Schema 不匹配。3. 连接参数如 URI、端口错误。1. 对于已存在的集合初始化Milvus时设置auto_schemaFalse。2. 检查集合已有字段确保metadata中的键能对应上或使用collection.create_partition管理。3. 使用connections.list_connections()检查连接状态。按条件删除 (delete_by_condition) 报错1. 表达式语法错误。2. 尝试在未加载的集合上执行删除。3. 字段名或值类型不匹配。1. 使用简单的表达式测试如pk id1。Milvus 表达式语法参考其文档。2. 确保collection.load()已被调用。3. 使用collection.query(expr, output_fields[*])先验证表达式是否能正确查询到目标数据。嵌入维度不匹配1. 创建集合时定义的dim与 Embedding 模型实际产生的维度不一致。1. 确认 Embedding 模型的输出维度如text-embedding-ada-002是 1536。2. 重建集合或使用collection.alter修改字段此操作复杂通常建议重建。性能问题插入/更新慢1. 单条插入。2. 索引类型不适合写多读少的场景。3. 网络延迟或资源不足。1. 采用批量插入如每次插入 100-500 条数据。2. 评估索引类型IVF_FLAT比HNSW写入更快但搜索精度可能稍低。3. 监控 Milvus 集群资源CPU、内存、磁盘IO。6. 最佳实践与工程建议将 LangChain 与 Milvus 用于生产环境时除了功能实现更应关注可靠性、可维护性和性能。主键设计是基石全局唯一且业务相关不要使用无意义的自增ID或UUID。使用能标识数据实体的字段如user_id:content_hash的组合便于直接定位和业务关联。类型选择根据数据量选择VARCHAR或INT64。VARCHAR更灵活INT64性能稍好。实现数据版本化管理在元数据中增加version、update_time字段。实施“软删除”增加is_deleted标志位而不是物理删除。这样便于数据追溯和恢复。重大更新时可以采用“插入新版本标记旧版本失效”的策略而非原地更新便于AB测试或回滚。批处理与异步操作对于大规模数据导入或更新务必使用批量操作。pymilvus的insert方法接受多行数据。考虑将耗时的 DML 操作放入异步任务队列如 Celery、Dramatiq避免阻塞主应用线程。# 批量插入示例 def batch_insert_docs(doc_list: List[Document], batch_size200): for i in range(0, len(doc_list), batch_size): batch doc_list[i:ibatch_size] # ... 处理并插入 batch self.collection.insert(batch_data) self.collection.flush()操作前备份与验证在执行批量删除或更新前务必先执行查询确认影响范围如示例中的delete_by_condition方法所示。对于核心数据定期为 Milvus 集合创建快照如果使用 Milvus 企业版或云服务。编写数据迁移脚本时应有“试运行”Dry-Run模式只打印将要执行的操作而不实际执行。监控与日志记录所有 DML 操作的关键信息操作类型、影响的主键、操作时间、执行人或服务、耗时。监控 Milvus 的关键指标集合中的实体数量、索引状态、查询/插入 QPS、内存使用率。为删除操作设置更高级别的日志如 WARNING 或 ERROR 级别。与 LangChain 生态的优雅结合将MilvusDMLManager这样的管理类封装成独立的服务或模块与 LangChain 的 Chain 或 Agent 通过 API 交互。利用 LangChain 的Retriever抽象。你可以自定义一个Retriever它在内部调用你的管理器并可能加入缓存、过滤等逻辑。考虑使用langchain.indexes模块来维护向量存储的增量更新它提供了SQLRecordManager来追踪文档的哈希状态。通过以上实践你可以构建一个不仅功能强大而且稳定、可观测、易于维护的向量数据管理系统使其成为你 AI 应用坚实的数据底座。