从0开始进阶AI Agent--AI智能体开发工程师 篇章4

——— 一个完整的RAG智能体应用小工程

先来看一下该工程的结构:

AiAgent/ ├── app_complete.py # 主程序 ├── db_manager.py # 数据库管理模块 ├── .env # API密钥 ├── knowledge_base/ # 文档存放目录 │ ├── sample.txt │ └── ... └── chat_history.db # 自动生成的数据库文件

这次废话不多说,直接上代码干货我们再开始讲述,先贴上项目的运行结果图。

1. 数据库管理模块(db_manager.py

# db_manager.py import sqlite3 import json from datetime import datetime from langchain_core.messages import HumanMessage, AIMessage, SystemMessage DB_PATH = "chat_history.db" def init_db(): """初始化数据库,创建必要的表""" conn = sqlite3.connect(DB_PATH) c = conn.cursor() # 创建对话历史表 c.execute('''CREATE TABLE IF NOT EXISTS chat_history (id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT, role TEXT, content TEXT, timestamp DATETIME DEFAULT CURRENT_TIMESTAMP)''') # 创建用户偏好表(用于长期记忆) c.execute('''CREATE TABLE IF NOT EXISTS user_preferences (id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT UNIQUE, preferences TEXT, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP)''') conn.commit() conn.close() def save_message(session_id: str, role: str, content: str): """保存一条消息到数据库""" conn = sqlite3.connect(DB_PATH) c = conn.cursor() c.execute("INSERT INTO chat_history (session_id, role, content) VALUES (?, ?, ?)", (session_id, role, content)) conn.commit() conn.close() def load_history(session_id: str, limit: int = 20): """加载指定会话的最近N条消息(按时间正序)""" conn = sqlite3.connect(DB_PATH) c = conn.cursor() c.execute("SELECT role, content FROM chat_history WHERE session_id=? ORDER BY timestamp DESC LIMIT ?", (session_id, limit * 2)) # limit轮数 rows = c.fetchall() conn.close() messages = [] for role, content in reversed(rows): # 按时间正序 if role == "user": messages.append(HumanMessage(content=content)) elif role == "assistant": messages.append(AIMessage(content=content)) elif role == "system": messages.append(SystemMessage(content=content)) return messages def save_preference(session_id: str, preferences: dict): """保存用户偏好(长期记忆)""" conn = sqlite3.connect(DB_PATH) c = conn.cursor() c.execute("INSERT OR REPLACE INTO user_preferences (session_id, preferences) VALUES (?, ?)", (session_id, json.dumps(preferences, ensure_ascii=False))) conn.commit() conn.close() def load_preference(session_id: str) -> dict: """加载用户偏好""" conn = sqlite3.connect(DB_PATH) c = conn.cursor() c.execute("SELECT preferences FROM user_preferences WHERE session_id=?", (session_id,)) row = c.fetchone() conn.close() if row: return json.loads(row[0]) return {} # 初始化数据库 init_db() print("✅ 数据库初始化完成")

2. 多智能体协作模块(在app_complete.py中实现)

创建三个智能体角色:

  • 研究员(Researcher):负责从知识库检索信息

  • 分析师(Analyst):负责分析信息并提取关键点

  • 写作助手(Writer):负责组织语言并生成最终回答

3. 主程序(app_complete.py

# app_complete.py import streamlit as st import os import warnings import asyncio from typing import TypedDict, List, Dict, Any from datetime import datetime from dotenv import load_dotenv warnings.filterwarnings("ignore") from langchain_openai import ChatOpenAI from langchain_community.embeddings import HuggingFaceEmbeddings from langchain_community.document_loaders import TextLoader, PyPDFLoader, Docx2txtLoader from langchain_text_splitters import RecursiveCharacterTextSplitter from langchain_community.vectorstores import FAISS from langchain_core.tools import tool from langchain.agents import create_agent from langchain_core.messages import SystemMessage, HumanMessage, AIMessage from langgraph.graph import StateGraph, END from langgraph.checkpoint import MemorySaver from db_manager import save_message, load_history, save_preference, load_preference load_dotenv() # ---------- 状态定义(用于多智能体) ---------- class MultiAgentState(TypedDict): question: str session_id: str retrieved_docs: List[str] analysis: str final_answer: str user_preferences: Dict[str, Any] # ---------- 加载向量库 ---------- @st.cache_resource def load_vector_store(): doc_dir = "./knowledge_base" if not os.path.exists(doc_dir): os.makedirs(doc_dir) st.error(f"📁 已创建 {doc_dir} 文件夹,请放入文档后重启") return None documents = [] for file in os.listdir(doc_dir): file_path = os.path.join(doc_dir, file) try: if file.endswith(".txt"): loader = TextLoader(file_path, encoding="utf-8") elif file.endswith(".pdf"): loader = PyPDFLoader(file_path) elif file.endswith(".docx"): loader = Docx2txtLoader(file_path) else: continue docs = loader.load() if docs: documents.extend(docs) print(f"✅ 已加载: {file}") except Exception as e: print(f"⚠️ 跳过 {file}: {e}") if not documents: return None text_splitter = RecursiveCharacterTextSplitter( chunk_size=500, chunk_overlap=50, separators=["\n\n", "\n", "。", "!", "?", ";", ",", " ", ""], ) chunks = text_splitter.split_documents(documents) embeddings = HuggingFaceEmbeddings( model_name="paraphrase-multilingual-MiniLM-L12-v2", model_kwargs={'device': 'cpu'}, encode_kwargs={'normalize_embeddings': True} ) vector_store = FAISS.from_documents(chunks, embeddings) return vector_store # ---------- 创建检索工具 ---------- def create_retriever_tool(vector_store): retriever = vector_store.as_retriever(search_kwargs={"k": 5}) @tool def search_knowledge(query: str) -> str: """从本地知识库中检索与用户问题最相关的信息片段。""" docs = retriever.invoke(query) if not docs: return "未找到相关信息。" return "\n\n---\n\n".join([doc.page_content for doc in docs]) return search_knowledge # ---------- 创建多智能体协作图 ---------- def create_multi_agent_system(vector_store, llm): """构建一个包含研究员、分析师、写作助手的多智能体协作系统""" retriever_tool = create_retriever_tool(vector_store) # 节点1:研究员 - 负责检索信息 def researcher_node(state: MultiAgentState) -> Dict: """研究员:从知识库检索相关信息""" question = state["question"] # 如果用户有偏好,调整检索策略 prefs = state.get("user_preferences", {}) if prefs.get("prefer_detailed", False): # 调大检索数量 docs = vector_store.similarity_search(question, k=8) else: docs = vector_store.similarity_search(question, k=5) retrieved_texts = [doc.page_content for doc in docs] return { "retrieved_docs": retrieved_texts, "analysis": f"检索到 {len(retrieved_texts)} 条相关信息" } # 节点2:分析师 - 负责提取关键信息 def analyst_node(state: MultiAgentState) -> Dict: """分析师:提取关键信息,形成分析结论""" docs_text = "\n\n".join(state.get("retrieved_docs", [""])[:3]) # 取前3条 question = state["question"] prompt = f""" 你是一个数据分析师。请从以下文档中提取与用户问题最相关的关键信息。 用户问题:{question} 文档内容: {docs_text} 请用简洁的语言总结: 1. 核心答案(如果有直接答案) 2. 支持性证据(3-5个关键点) 3. 不确定或需要补充的地方 """ response = llm.invoke(prompt) return {"analysis": response.content} # 节点3:写作助手 - 负责生成最终回答 def writer_node(state: MultiAgentState) -> Dict: """写作助手:基于分析和用户偏好组织最终回答""" analysis = state.get("analysis", "暂无分析结果") question = state["question"] prefs = state.get("user_preferences", {}) # 根据用户偏好调整风格 style_note = "" if prefs.get("prefer_concise", False): style_note = "请用简洁的语言回答,控制在100字以内。" elif prefs.get("prefer_detailed", False): style_note = "请提供详细的回答,充分解释每个要点。" prompt = f""" 你是一个专业的知识写作助手。请根据以下分析,生成对用户问题的最终回答。 用户问题:{question} 分析结果: {analysis} 写作要求:保持专业、准确、易懂。 {style_note} """ response = llm.invoke(prompt) return {"final_answer": response.content} # 构建图 builder = StateGraph(MultiAgentState) builder.add_node("researcher", researcher_node) builder.add_node("analyst", analyst_node) builder.add_node("writer", writer_node) builder.set_entry_point("researcher") builder.add_edge("researcher", "analyst") builder.add_edge("analyst", "writer") builder.add_edge("writer", END) # 编译并添加记忆(使用MemorySaver实现短期记忆) memory = MemorySaver() graph = builder.compile(checkpointer=memory) return graph # ---------- 创建简单智能体(用于对比) ---------- def create_simple_agent(vector_store): """创建单智能体(用于对比)""" retriever_tool = create_retriever_tool(vector_store) llm = ChatOpenAI( model="deepseek-v4-pro", # 替换为你实际的模型名 base_url="https://api.deepseek.com/v1", temperature=0, ) agent = create_agent( model=llm, tools=[retriever_tool], system_prompt="你是一个知识渊博的助手,基于知识库回答问题,结合历史对话保持连贯。" ) return agent # ---------- Streamlit 主界面 ---------- st.set_page_config(page_title="RAG 智能体 - 完整版", page_icon="🤖", layout="wide") # 侧边栏配置 with st.sidebar: st.title("⚙️ 配置") # 选择工作模式 mode = st.radio( "选择智能体模式", ["🚀 多智能体协作(推荐)", "🤖 单智能体(传统)"], help="多智能体协作模式会使用三个智能体分工合作,回答质量更高" ) st.divider() # 用户偏好设置(长期记忆) st.subheader("🧠 记忆设置") prefer_concise = st.checkbox("偏好简洁回答", value=False) prefer_detailed = st.checkbox("偏好详细回答", value=True) if st.button("💾 保存偏好设置"): save_preference("default_session", { "prefer_concise": prefer_concise, "prefer_detailed": prefer_detailed, "updated_at": datetime.now().isoformat() }) st.success("✅ 偏好已保存!") st.divider() # 显示历史统计 st.subheader("📊 统计信息") history = load_history("default_session", limit=100) st.metric("历史对话轮数", len(history)//2) # 主界面 st.title("🤖 RAG 智能体 - 完整版") st.caption("支持多智能体协作 · 跨会话记忆 · 个性化偏好") # 加载向量库 vector_store = load_vector_store() if vector_store is None: st.error("❌ 请将文档放入 knowledge_base 目录后重启应用") st.stop() # 初始化会话状态 if "messages" not in st.session_state: st.session_state.messages = [] # 从数据库加载历史 db_history = load_history("default_session", limit=20) for msg in db_history: if isinstance(msg, HumanMessage): st.session_state.messages.append({"role": "user", "content": msg.content}) elif isinstance(msg, AIMessage): st.session_state.messages.append({"role": "assistant", "content": msg.content}) # 加载用户偏好 user_prefs = load_preference("default_session") # 初始化LLM llm = ChatOpenAI( model="deepseek-v4-pro", # 替换为你实际的模型名 base_url="https://api.deepseek.com/v1", temperature=0, ) # 根据模式创建不同的智能体 if mode == "🚀 多智能体协作(推荐)": if "multi_agent_graph" not in st.session_state: st.session_state.multi_agent_graph = create_multi_agent_system(vector_store, llm) agent_obj = st.session_state.multi_agent_graph agent_type = "multi" else: if "simple_agent" not in st.session_state: st.session_state.simple_agent = create_simple_agent(vector_store) agent_obj = st.session_state.simple_agent agent_type = "single" # 显示历史消息 for msg in st.session_state.messages: with st.chat_message(msg["role"]): st.markdown(msg["content"]) # 输入框 if prompt := st.chat_input("请输入你的问题"): # 显示用户消息 with st.chat_message("user"): st.markdown(prompt) st.session_state.messages.append({"role": "user", "content": prompt}) # 保存到数据库 save_message("default_session", "user", prompt) # 调用智能体 with st.chat_message("assistant"): with st.spinner("🤔 思考中..."): try: if agent_type == "multi": # 多智能体调用 config = {"configurable": {"thread_id": "default_session"}} state_input = { "question": prompt, "session_id": "default_session", "user_preferences": user_prefs, "retrieved_docs": [], "analysis": "", "final_answer": "" } result = agent_obj.invoke(state_input, config) answer = result.get("final_answer", "未生成回答") else: # 单智能体调用 history_msgs = load_history("default_session", limit=10) full_messages = [ SystemMessage(content="你是一个知识渊博的助手,基于知识库回答问题。") ] full_messages.extend(history_msgs) full_messages.append(HumanMessage(content=prompt)) result = agent_obj.invoke({"messages": full_messages}) answer = result["messages"][-1].content st.markdown(answer) except Exception as e: st.error(f"❌ 发生错误: {e}") answer = f"抱歉,发生了错误:{e}" # 保存回答 st.session_state.messages.append({"role": "assistant", "content": answer}) save_message("default_session", "assistant", answer) # 自动学习用户偏好(根据提问内容更新) if "详细" in prompt or "深入" in prompt: new_prefs = user_prefs.copy() new_prefs["prefer_detailed"] = True save_preference("default_session", new_prefs) st.toast("🧠 已自动学习你的偏好:详细回答") # 底部信息 st.divider() st.caption("💡 提示:多智能体模式会使用研究员→分析师→写作助手三个角色协作回答")

4. 安装依赖

确保aiagent环境下安装了所有必要的包:

pip install streamlit langchain langchain-openai langchain-community langchain-text-splitters faiss-cpu sentence-transformers langgraph sqlite3

可用清华源进行加速:

pip install streamlit langchain langchain-openai langchain-community langchain-text-splitters faiss-cpu sentence-transformers langgraph sqlite3 -i https://pypi.tuna.tsinghua.edu.cn/simp

5. 运行

streamlit run app_complete.py

🎯 本示例包含了什么

功能实现方式位置
Web界面Streamlit主界面
长期记忆SQLite +db_manager.py跨会话保存对话
内置记忆组件LangGraph的MemorySaver多智能体内部记忆
多智能体协作LangGraph状态图研究员→分析师→写作助手
用户偏好学习SQLite存储偏好自动学习用户风格偏好
文档检索FAISS + HuggingFace Embedding本地RAG

运行效果

  • 左侧边栏可以切换单智能体/多智能体模式,方便对比效果。

  • 你可以保存用户偏好(简洁/详细),系统会记住并在后续回答中应用。

  • 所有对话都会自动存入SQLite,关闭程序再打开,历史依然存在。

  • 多智能体模式下,系统会先检索,再分析,最后生成回答,质量更高

🎉 现在可以在浏览器中访问http://localhost:8501查看完整RAG智能体应用了。

🚀现在能用它做什么?

1. 测试基础功能

  • 上传文档:在knowledge_base文件夹中放入一些文档(.txt, .pdf, .docx),重启应用后会自动加载。

  • 开始对话:在输入框中提问,智能体会基于文档回答问题。

  • 切换模式:在左侧边栏切换"单智能体"和"多智能体协作"模式,对比回答效果。

2. 测试记忆功能

  • 关闭浏览器标签页,重新打开http://localhost:8501,对话历史会自动加载(SQLite持久化)。

  • 在左侧边栏设置"偏好简洁回答"或"偏好详细回答",系统会记住个人偏好。

3. 测试多智能体协作

  • 选择"多智能体协作(推荐)"模式

  • 问一个复杂问题,比如:"请详细介绍RAG技术的原理和应用"

  • 系统会依次调用:研究员(检索)→ 分析师(分析)→ 写作助手(生成回答)

🔧 常见问题及解决方案

问题1:打开页面后报错 "未找到任何有效文档"

解决:在项目根目录下创建knowledge_base文件夹,放入至少一个文档文件(.txt, .pdf, .docx),然后刷新页面。

问题2:模型API调用失败(401/404)

解决:检查.env文件中的OPENAI_API_KEY是否正确,以及modelbase_url是否与你的API服务匹配。

问题3:Streamlit界面卡顿或加载慢

原因:首次运行会下载HuggingFace模型(约470MB),需要耐心等待。
解决:可以改用更小的模型,在load_vector_store函数中替换:

embeddings = HuggingFaceEmbeddings( model_name="distiluse-base-multilingual-cased-v1", # 更小的模型,约250MB model_kwargs={'device': 'cpu'} )

📊 功能对比表

功能单智能体模式多智能体模式
响应速度⚡ 快🐢 稍慢(多步推理)
回答质量⭐⭐⭐ 好⭐⭐⭐⭐⭐ 优秀
处理复杂问题⭐⭐ 一般⭐⭐⭐⭐⭐ 擅长
资源消耗💚 低💛 中高

🎯 下一步可以做的

  1. 定制自己的多智能体:修改create_multi_agent_system中的节点(researcher/analyst/writer),调整它们的提示词或增加新角色(比如"校对员"、"翻译官"等)。

  2. 增加更多用户偏好:在save_preference中增加新字段,比如"回答语言偏好(中文/英文)"、"行业术语偏好"等。

  3. 部署到云端:使用Streamlit Cloud、HuggingFace Spaces或自己的服务器部署应用,让其他人也能使用。

  4. 增加文件上传功能:在Streamlit界面中添加文件上传组件,让用户直接在浏览器中上传文档,无需手动放到knowledge_base文件夹。

💡 更加专业的改进

# 在 app_complete.py 顶部添加以下代码,让Streamlit支持文件上传 uploaded_files = st.sidebar.file_uploader( "📤 上传知识库文档", type=["txt", "pdf", "docx"], accept_multiple_files=True ) if uploaded_files: for file in uploaded_files: # 保存文件到 knowledge_base 目录 with open(os.path.join("knowledge_base", file.name), "wb") as f: f.write(file.getbuffer()) st.sidebar.success("✅ 文档上传成功!请重启应用以加载新文档。")