Chapter3 Milvus 在 AI Agent 中的应用 ipynb可执行代码请点击:Milvus在AI Agent中的应用.ipynb 本项目将通过LangGraph和Milvus构建Agent 什么是Agent 一个典型的Agent通常包含以下几个核心组件: palnning(规划): 目标设定和分解:Agent首先需要理解用户的最终目标,并将其分解为一些了可执行的子任务或步骤 策略选择:对于每个子任务,Agent可能有多种执行方式或则工具可以选择。规划模块负责选择最优策略。 执行监控和调整:在任务执行过程中,Agent需要监控进展,并根据实际情况调整计划。
ipynb可执行代码请点击:Milvus在AI Agent中的应用.ipynb
本项目将通过LangGraph和Milvus构建Agent
一个典型的Agent通常包含以下几个核心组件:
Milvus 作为一个高性能的向量数据库,在 AI Agent 中可以扮演两个至关重要的角色:
通过将信息向量化并存储在 Milvus 中,Agent 可以利用语义相似性来检索知识和记忆,而不仅仅是关键词匹配,这使得 Agent 的信息获取和利用能力大大增强。
接下来,我们将通过一个简化的案例,演示一个 AI Agent 如何利用 Milvus 作为外部知识库。我们将使用 LangGraph 来构建 Agent 的控制流程。
核心流程:
同时,我们也可以构想 Agent 如何存储对话片段:
下面我们聚焦于使用 Milvus 作为外部知识库的 LangGraph Agent 实现。
我们将简化上述参考链接中的 GraphRAG 概念,构建一个更直接的 Agent,它有一个工具是查询 Milvus。
! pip install pymilvus langchain==0.3.25 langgraph==0.4.7 langchain_openai==0.3.18 langchain_community==0.0.38 langchain-core==0.3.61
import os import uuid from typing import TypedDict, Annotated, List, Union import operator from datetime import datetime from langchain_community.chat_models import ChatZhipuAI from langchain_community.embeddings import ZhipuAIEmbeddings from langchain_core.messages import BaseMessage, HumanMessage, AIMessage, ToolMessage from langchain_core.tools import tool from langchain_community.vectorstores import Milvus from langgraph.graph import StateGraph, END from langgraph.prebuilt import create_react_agent from pymilvus import connections, utility, CollectionSchema, FieldSchema, DataType, Collection # 中文场景下,智谱的效果好一些,所以这里记得填写密钥 ZHIPUAI_API_KEY = "" # 记得启动你本地的Milvus服务 MILVUS_HOST = "localhost" MILVUS_PORT = "19530" MILVUS_COLLECTION_NAME = "ai_agent_knowledge_base" MILVUS_EMBEDDING_DIM = 1024 ID_FIELD_NAME = "doc_id" TEXT_FIELD_NAME = "text_content" VECTOR_FIELD_NAME = "embedding" llm = ChatZhipuAI( model="glm-4", temperature=0, api_key=ZHIPUAI_API_KEY ) embeddings_model = ZhipuAIEmbeddings( model="embedding-2", api_key=ZHIPUAI_API_KEY ) print("配置加载完毕(使用 GLM + Milvus)。") # 初始化向量数据库连接以及字段shema和collection def init_milvus_collection(): connections.connect(host=MILVUS_HOST, port=MILVUS_PORT) if utility.has_collection(MILVUS_COLLECTION_NAME): print(f"集合 {MILVUS_COLLECTION_NAME} 已存在。") return doc_id = FieldSchema(name=ID_FIELD_NAME, dtype=DataType.VARCHAR, is_primary=True, max_length=36) text_content = FieldSchema(name=TEXT_FIELD_NAME, dtype=DataType.VARCHAR, max_length=65535) embedding = FieldSchema(name=VECTOR_FIELD_NAME, dtype=DataType.FLOAT_VECTOR, dim=MILVUS_EMBEDDING_DIM) schema = CollectionSchema(fields=[doc_id, text_content, embedding], description="AI Agent Knowledge Base") collection = Collection(name=MILVUS_COLLECTION_NAME, schema=schema) index_params = {"index_type": "HNSW", "metric_type": "COSINE", "params": {"M": 8, "efConstruction": 64}} collection.create_index(VECTOR_FIELD_NAME, index_params) print(f"集合 {MILVUS_COLLECTION_NAME} 创建成功。") # langgraph组件 vectorstore = Milvus( embedding_function=embeddings_model, connection_args={"host": MILVUS_HOST, "port": MILVUS_PORT}, collection_name=MILVUS_COLLECTION_NAME, auto_id=False, primary_field=ID_FIELD_NAME, text_field=TEXT_FIELD_NAME, vector_field=VECTOR_FIELD_NAME, ) # 从向量数据库创建一个检索器,并配置每次检索时返回最相似的前三条结果 retriever = vectorstore.as_retriever(search_kwargs={"k": 3}) # 定义工具,可不要随便起abc的函数名称! @tool def search_knowledge(query: str) -> str: """从 Milvus 知识库中检索相关信息""" docs = retriever.invoke(query) if not docs: return "未找到相关信息。" return "\n".join([f"[{i+1}] {doc.page_content}" for i, doc in enumerate(docs)]) @tool def get_current_time(placeholder: str = "default") -> str: """Returns the current date and time.""" print("\n[Tool Call: get_current_time]") return datetime.now().strftime("%Y-%m-%d %H:%M:%S") # 工具列表 tools = [search_knowledge,get_current_time] agent_executor = create_react_agent(llm, tools) # 你可以尝试询问当前几点了,验证工具是否起作用 if __name__ == "__main__": print("\n GLM + Milvus 智能体已就绪!输入问题开始对话(输入 'quit' 退出):\n") user_input = input(" 你: ").strip() try: response = agent_executor.invoke({"messages": [("human", user_input)]}) ai_message = response["messages"][-1].content print(f" AI: {ai_message}\n") except Exception as e: print(f"❌ 错误: {e}\n")
# Milvus Setup and Helper Functions def connect_to_milvus(): """建立与 Milvus 的连接""" try: connections.connect(host=MILVUS_HOST, port=MILVUS_PORT) print(f"成功连接到 Milvus: {MILVUS_HOST}:{MILVUS_PORT}") except Exception as e: print(f"连接 Milvus 失败: {e}") raise def create_milvus_collection_if_not_exists(): """如果集合不存在,则创建它""" connect_to_milvus() # 确保连接 if utility.has_collection(MILVUS_COLLECTION_NAME): print(f"集合 '{MILVUS_COLLECTION_NAME}' 已存在.") utility.drop_collection(collection_name=MILVUS_COLLECTION_NAME) field_id = FieldSchema(name=ID_FIELD_NAME, dtype=DataType.VARCHAR, is_primary=True, max_length=36) field_text = FieldSchema(name=TEXT_FIELD_NAME, dtype=DataType.VARCHAR, max_length=65535) # 存储原始文本 field_embedding = FieldSchema(name=VECTOR_FIELD_NAME, dtype=DataType.FLOAT_VECTOR, dim=MILVUS_EMBEDDING_DIM) schema = CollectionSchema( fields=[field_id, field_text, field_embedding], description="AI Agent Knowledge Base collection", enable_dynamic_field=False # 动态字段 如果需要额外元数据且不想预定义,可以设为True ) collection = Collection(MILVUS_COLLECTION_NAME, schema=schema) print(f"集合 '{MILVUS_COLLECTION_NAME}' 创建成功.") # 为向量字段创建索引 index_params = { "metric_type": "L2", # 或 "IP" "index_type": "IVF_FLAT", "params": {"nlist": 128}, } collection.create_index(field_name=VECTOR_FIELD_NAME, index_params=index_params) print(f"为字段 '{VECTOR_FIELD_NAME}' 创建索引成功.") collection.load() print(f"集合 '{MILVUS_COLLECTION_NAME}' 已加载.") return collection def insert_data_to_milvus(collection: Collection, texts: List[str]): """将文本数据向量化并插入 Milvus""" if not texts: print("没有数据需要插入。") return print(f"正在为 {len(texts)} 条文本生成向量...") vectors = embeddings_model.embed_documents(texts) print("向量生成完毕。") # 准备插入数据 data_to_insert = [] for i, text_content in enumerate(texts): data_to_insert.append({ ID_FIELD_NAME: str(uuid.uuid4()), TEXT_FIELD_NAME: text_content, VECTOR_FIELD_NAME: vectors[i] }) print(f"正在向 Milvus 集合 '{collection.name}' 插入 {len(data_to_insert)} 条数据...") insert_result = collection.insert(data_to_insert) collection.flush() # 确保数据持久化 print(f"数据插入成功. 影响行数: {insert_result.insert_count}") print(f"当前集合实体数量: {collection.num_entities}") # 执行 Milvus 初始化 try: knowledge_collection = create_milvus_collection_if_not_exists() # 准备一些示例知识数据 (仅在首次运行时或需要时插入) # 为避免重复插入,可以检查集合是否为空 if knowledge_collection.num_entities == 0: print("知识库为空,准备插入示例数据...") sample_knowledge = [ "Milvus 是一款开源的向量数据库,专为大规模向量相似性搜索和分析而设计。", "AI Agent 可以利用 Milvus 作为其长期记忆存储和外部知识库。", "LangGraph 是一个用于构建有状态、多参与者应用程序的库,特别适合构建复杂的 AI Agent。", "向量数据库通过将数据转换为向量嵌入,并使用专门的索引进行高效的相似性搜索。", "RAG (Retrieval Augmented Generation) 是一种结合了检索系统和生成模型的AI技术,可以提高生成内容的准确性和相关性。", "太阳是太阳系的中心天体,其核心温度高达1500万摄氏度。", "Python 是一种广泛使用的高级编程语言,以其简洁的语法和强大的库生态系统而闻名。", "DataWhale 是国内领先的 AI 开源学习社区,成立于 2018 年,致力于推动人工智能领域的开源教育与协作学习。", "DataWhale 社区覆盖全球 3500 多所高校,拥有数百万开发者,所有学习资料和项目代码均开源在 GitHub 上。", "DataWhale 通过组织黑客松、组队学习和开源项目共建,帮助开发者系统掌握机器学习、大模型、向量数据库等前沿技术。", "DataWhale 与 AMD、魔搭社区等机构合作,共同推动 ROCm 生态和国产 AI 基础设施的开发者生态建设。" ] insert_data_to_milvus(knowledge_collection, sample_knowledge) else: print(f"知识库中已有 {knowledge_collection.num_entities} 条数据,跳过示例数据插入。") except Exception as e: print(f"Milvus 初始化或数据插入过程中发生错误: {e}") # 在Notebook中,我们可能不希望程序因Milvus连接问题而完全停止后续单元格的执行 # 但后续依赖Milvus的单元格可能会失败 knowledge_collection = None # 标记为None,以便后续检查
from typing import List, TypedDict, Annotated import operator from datetime import datetime from langgraph.graph import StateGraph, END from langchain_core.tools import tool from langchain_core.messages import BaseMessage, ToolMessage from langchain_core.runnables import RunnableLambda from langchain_core.utils.function_calling import convert_to_openai_tool # 1. 定义工具 @tool def search_milvus_knowledge_base(query: str) -> str: """ 从Milvus中查找与问题相关的信息。 输入的问题应是一个特定于milvus中存储的数据的问题 """ if not knowledge_collection: return "Milvus knowledge base is not available." print(f"\n[Tool Call: search_milvus_knowledge_base] Query: {query}") query_vector = embeddings_model.embed_query(query) search_params = { "metric_type": "L2", "params": {"nprobe": 10}, # 基于索引类型和数据规模来定义nprobe大小 } # 执行搜索 results = knowledge_collection.search( data=[query_vector], anns_field=VECTOR_FIELD_NAME, param=search_params, limit=3, # 返回前三条最相关的答案 expr=None, # 可选择的标量过滤语句 output_fields=[TEXT_FIELD_NAME] # 返回原始的内容字段 ) context = "" if results and results[0]: context_docs = [hit.entity.get(TEXT_FIELD_NAME) for hit in results[0]] context = "\n".join(context_docs) print(f"[Tool Result] Found context: {context[:200]}...") else: print("[Tool Result] No relevant context found in Milvus.") context = "No relevant information found in the knowledge base." return context @tool def get_current_time(placeholder: str = "default") -> str: # Langchain tools often expect an input arg """Returns the current date and time.""" print("\n[Tool Call: get_current_time]") return datetime.now().strftime("%Y-%m-%d %H:%M:%S") # 定义工具列表 tools = [search_milvus_knowledge_base,get_current_time] # 2. 定义Agent状态 # 这是典型的对话历史的存储方式 class AgentState(TypedDict): messages: Annotated[List[BaseMessage], operator.add] # 3. 定义节点 def agent_node(state: AgentState) -> dict: """ Agent node: 决定下一个行动是什么 (是调用工具还是直接回复). """ print("\n[Node: Agent]") # 向llm展示可用的工具都有哪些 bound_llm = llm.bind_tools(tools) response = bound_llm.invoke(state["messages"]) # 展示Agent的决定 print(f"[Agent Decision] Response: {response.content}, Tool Calls: {response.tool_calls}") return {"messages": [response]} def tool_node(state: AgentState) -> dict: """ Tool node: 执行agent发起的工具调用 """ print("\n[Node: Tool Executor]") last_message = state["messages"][-1] if not hasattr(last_message, "tool_calls") or not last_message.tool_calls: print("[Tool Executor] No tool calls found in the last message.") return {"messages": []} tool_messages = [] for tool_call in last_message.tool_calls: tool_name = tool_call["name"] tool_input = tool_call["args"] # 通过name查找对应的tool tool = next((t for t in tools if t.name == tool_name), None) if not tool: tool_messages.append( ToolMessage( content=f"Error: Tool {tool_name} not found.", tool_call_id=tool_call["id"] ) ) continue try: # 执行工具 result = tool.invoke(tool_input) tool_messages.append( ToolMessage( content=str(result), tool_call_id=tool_call["id"] ) ) except Exception as e: tool_messages.append( ToolMessage( content=f"Error executing tool {tool_name}: {str(e)}", tool_call_id=tool_call["id"] ) ) print(f"[Tool Executor] Executed tools, results: {tool_messages}") return {"messages": tool_messages} # 4. 定义情景上的边界 def should_continue(state: AgentState) -> str: """ 决定什么时候继续调用工具还是结束 """ print("\n[Edge: should_continue]") last_message = state["messages"][-1] if hasattr(last_message, "tool_calls") and last_message.tool_calls: print("[Edge Decision] Continue to 'tools'") return "tools" print("[Edge Decision] End") return END # 5. 构建图 workflow = StateGraph(AgentState) # 添加节点 workflow.add_node("agent", RunnableLambda(agent_node)) workflow.add_node("tools", RunnableLambda(tool_node)) # 设置入口 workflow.set_entry_point("agent") # 添加边界 workflow.add_conditional_edges( "agent", should_continue, { "tools": "tools", END: END } ) # 添加工具到agent的连接 workflow.add_edge("tools", "agent") # 编译 app = workflow.compile() print("\nLangGraph App compiled successfully!")
from IPython.display import Image, display try: display(Image(app.get_graph().draw_mermaid_png())) except Exception: pass
if not knowledge_collection: print(" Milvus 连接和集合初始化失败。") else: print("Agent 已准备就绪。开始提问吧!(输入 'exit' 退出)") print("-" * 30) # 演示 Agent 使用 Milvus 作为外部知识库 print("\n--- 案例1: Agent 利用 Milvus 查找信息 ---") query1 = "Milvus 是什么?" print(f"User: {query1}") inputs = {"messages": [HumanMessage(content=query1)]} # 使用 stream 方法逐步查看执行过程 for event in app.stream(inputs): for key, value in event.items(): print(f"--- Event for Node: {key} ---") if "messages" in value: # 打印最新消息的内容 latest_message = value["messages"][-1] if isinstance(latest_message, AIMessage): print(f"AI: {latest_message.content}") if latest_message.tool_calls: print(f"AI requests tool call: {latest_message.tool_calls}") elif isinstance(latest_message, ToolMessage): print(f"Tool Result ({latest_message.tool_call_id}): {latest_message.content}") else: print(f"Message ({type(latest_message).__name__}): {latest_message.content}") print("-" * 10) print("\n--- 案例2: Agent 回答一个不需要查知识库的问题 (可能直接回答或拒绝) ---") query2 = "你好吗?" print(f"User: {query2}") inputs = {"messages": [HumanMessage(content=query2)]} # 获取最终结果 final_response = app.invoke(inputs) if final_response and "messages" in final_response and final_response["messages"]: print(f"AI: {final_response['messages'][-1].content}") else: print("AI 未能生成回复。") # 演示 Agent 存储对话片段向量 (概念性,实际存储逻辑需要更完善) # 假设 query1 和其最终回复是一个需要记忆的片段 if final_response and "messages" in final_response: # 使用上一个交互的结果 print("\n--- 概念演示: 存储对话到 Milvus (作为记忆) ---") # 假设我们想要将用户的问题和Agent的最终回答作为一个记忆单元 # 这里的 final_response['messages'] 可能包含整个对话历史 # 我们通常取最后的用户问题和AI回答对 # 找到 query1 对应的最终 AIMessage # 这是一个简化的查找,实际中可能需要更复杂的逻辑来配对问答 q1_final_answer = "" # 假设 app.invoke 返回的 messages 列表的最后一个是最终答案 if final_response['messages'] and isinstance(final_response['messages'][-1], AIMessage): q1_final_answer = final_response['messages'][-1].content # 取决于上一个 invoke 的内容 # 如果 query1 导致了工具调用,我们可能需要从 stream 中找到它的最终回答 # 为了简化,我们直接使用上面交互中打印的最终回答 # 真实场景下,我们会捕获 app.invoke(inputs1) 的最终 AIMessage # 假设我们已经有了 query1 和 agent_final_answer_to_query1 # 这里我们手动设置一个示例,因为上面app.invoke(inputs)的最终结果是针对query2的 # 如果要精确获取query1的最终回答,需要重新运行app.invoke针对query1 # 或者从 app.stream 的事件中提取 # 为了演示,我们假设第一个问题"Milvus是什么"的最终答案是 "Milvus是一个开源的向量数据库..." (由LLM结合搜索结果生成) # 实际上,这个答案会在 stream 的某个 AIMessage 中出现 # 这里我们模拟一下,因为直接从上面的 stream 中捕获最终答案有点复杂 # 理想情况下,我们会有一个明确的 "final_answer" 状态或消息类型 # 假设我们通过某种方式获取到了 query1 的最终AI回答 simulated_final_answer_to_query1 = "Milvus 是一款先进的开源向量数据库,非常适合AI应用中的大规模相似性搜索。它能帮助Agent快速从大量文档中找到相关信息。" # 这是一个模拟的最终回答 if simulated_final_answer_to_query1: memory_text = f"用户问: {query1}\nAgent答: {simulated_final_answer_to_query1}" print(f"准备将以下对话片段存入记忆库:\n{memory_text}") # 为了避免与知识库冲突,可以存入不同的集合或使用分区 # 这里简单演示存入同一个集合,实际应用中应分开 try: # 为简化,我们假设有一个单独的记忆集合 memory_collection # memory_collection = create_milvus_collection_if_not_exists("ai_agent_memory", ...) # insert_data_to_milvus(memory_collection, [memory_text]) # 由于我们这里只有一个集合,就直接插入到 knowledge_collection,并作说明 print("注意: 实际应用中,对话记忆应存入专用集合或分区。此处为演示,插入当前知识库。") insert_data_to_milvus(knowledge_collection, [memory_text]) print("对话片段已(概念性地)存入 Milvus 记忆库。") # 如何在新对话开始时搜索相似历史 -> 召回相关记忆 new_user_query = "介绍一下向量数据库" # 一个新的,但与之前记忆相关的问题 print(f"\n新用户查询: {new_user_query}") print("Agent (概念上) 将搜索 Milvus 记忆库以查找相似历史对话...") # 实际操作: # 1. new_user_query_vector = embeddings_model.embed_query(new_user_query) # 2. search memory_collection with new_user_query_vector # 3. retrieved_memories = results_from_memory_collection # 4. Agent 使用 retrieved_memories 作为上下文辅助当前对话 # 这里我们用知识库搜索来模拟这个过程: retrieved_memories = search_milvus_knowledge_base(new_user_query) print(f"从Milvus中召回的(模拟的)相关记忆/知识:\n{retrieved_memories}") except Exception as e: print(f"存储或检索记忆时发生错误: {e}") else: print("未能获取到 query1 的最终回答,跳过记忆存储演示。")
Milvus 通过其强大的向量存储和检索能力,可以从多个方面赋能 AI Agent,使其更智能:
增强的知识获取与利用:
更强大的记忆能力:
提升任务执行效率与效果:
支持更复杂的 Agent 行为:
总之,Milvus 为 AI Agent 提供了一个坚实的数据基础,使其能够更有效地存储、管理和利用信息,从而在理解、规划、学习和交互等各个方面表现得更加智能。
目标: 体验并扩展我们刚刚构建的 AI Agent。
任务:
运行并理解 Agent:
search_milvus_knowledge_base 工具?app.stream(inputs) 的输出,理解 LangGraph 中节点的流转过程。扩展知识库:
Cell 3 (Milvus Setup and Helper Functions) 中,找到 sample_knowledge 列表。insert_data_to_milvus 函数。你可以:
utility.drop_collection(MILVUS_COLLECTION_NAME)(请谨慎操作!)。knowledge_collection.num_entities > 0,先 utility.drop_collection(MILVUS_COLLECTION_NAME),然后再调用 create_milvus_collection_if_not_exists() 和 insert_data_to_milvus()。请注意,这将删除所有现有数据。(可选) 尝试不同的查询:
(进阶可选) 添加一个新的简单工具:
get_current_time 工具,它不查询 Milvus,只是返回当前时间。
from datetime import datetime @tool def get_current_time(placeholder: str = "default") -> str: # Langchain tools often expect an input arg """Returns the current date and time.""" print("\n[Tool Call: get_current_time]") return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
Cell 4 的 tools 列表中: tools = [search_milvus_knowledge_base, get_current_time]。app = workflow.compile()。思考与记录:
Hands-on Exercise 3: 实操AI Agent Demo中的4个任务,相关答案在本章节对应的ipynb文件中,可以直接运行看到效果。
本文主要参考Milvus大佬的workshop项目,其中文字和代码部分全部来自该项目,但为了适配国内环境以及方便学习者使用,对代码部分进行了部分删改并为代码增加更详细的注释。