公司动态
从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(contentcontent)) elif role assistant: messages.append(AIMessage(contentcontent)) elif role system: messages.append(SystemMessage(contentcontent)) 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_asciiFalse))) 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, encodingutf-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_size500, chunk_overlap50, separators[\n\n, \n, 。, , , , , , ], ) chunks text_splitter.split_documents(documents) embeddings HuggingFaceEmbeddings( model_nameparaphrase-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, k8) else: docs vector_store.similarity_search(question, k5) 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(checkpointermemory) return graph # ---------- 创建简单智能体用于对比 ---------- def create_simple_agent(vector_store): 创建单智能体用于对比 retriever_tool create_retriever_tool(vector_store) llm ChatOpenAI( modeldeepseek-v4-pro, # 替换为你实际的模型名 base_urlhttps://api.deepseek.com/v1, temperature0, ) agent create_agent( modelllm, tools[retriever_tool], system_prompt你是一个知识渊博的助手基于知识库回答问题结合历史对话保持连贯。 ) return agent # ---------- Streamlit 主界面 ---------- st.set_page_config(page_titleRAG 智能体 - 完整版, page_icon, layoutwide) # 侧边栏配置 with st.sidebar: st.title(⚙️ 配置) # 选择工作模式 mode st.radio( 选择智能体模式, [ 多智能体协作推荐, 单智能体传统], help多智能体协作模式会使用三个智能体分工合作回答质量更高 ) st.divider() # 用户偏好设置长期记忆 st.subheader( 记忆设置) prefer_concise st.checkbox(偏好简洁回答, valueFalse) prefer_detailed st.checkbox(偏好详细回答, valueTrue) 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, limit100) 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, limit20) 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( modeldeepseek-v4-pro, # 替换为你实际的模型名 base_urlhttps://api.deepseek.com/v1, temperature0, ) # 根据模式创建不同的智能体 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, limit10) full_messages [ SystemMessage(content你是一个知识渊博的助手基于知识库回答问题。) ] full_messages.extend(history_msgs) full_messages.append(HumanMessage(contentprompt)) 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/simp5. 运行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是否正确以及model和base_url是否与你的API服务匹配。问题3Streamlit界面卡顿或加载慢原因首次运行会下载HuggingFace模型约470MB需要耐心等待。解决可以改用更小的模型在load_vector_store函数中替换embeddings HuggingFaceEmbeddings( model_namedistiluse-base-multilingual-cased-v1, # 更小的模型约250MB model_kwargs{device: cpu} ) 功能对比表功能单智能体模式多智能体模式响应速度⚡ 快 稍慢多步推理回答质量⭐⭐⭐ 好⭐⭐⭐⭐⭐ 优秀处理复杂问题⭐⭐ 一般⭐⭐⭐⭐⭐ 擅长资源消耗 低 中高 下一步可以做的定制自己的多智能体修改create_multi_agent_system中的节点researcher/analyst/writer调整它们的提示词或增加新角色比如校对员、翻译官等。增加更多用户偏好在save_preference中增加新字段比如回答语言偏好中文/英文、行业术语偏好等。部署到云端使用Streamlit Cloud、HuggingFace Spaces或自己的服务器部署应用让其他人也能使用。增加文件上传功能在Streamlit界面中添加文件上传组件让用户直接在浏览器中上传文档无需手动放到knowledge_base文件夹。 更加专业的改进# 在 app_complete.py 顶部添加以下代码让Streamlit支持文件上传 uploaded_files st.sidebar.file_uploader( 上传知识库文档, type[txt, pdf, docx], accept_multiple_filesTrue ) 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(✅ 文档上传成功请重启应用以加载新文档。)