场景 | 纯 RAG | 纯 Agent | RAG + Agent 混合 |
|---|---|---|---|
知识问答 | ✅ 基于文档检索 | ❌ 缺乏事实依据 | ✅ 检索增强,回答准确 |
多步推理 | ❌ 无法分解任务 | ✅ 可规划步骤 | ✅ 先检索知识,再规划行动 |
工具调用 | ❌ 无法执行操作 | ✅ 可调用 API | ✅ 根据上下文决策调用工具 |
记忆管理 | ❌ 无状态 | ✅ 支持短期记忆 | ✅ 结合外部知识长期记忆 |
企业智能运维助手典型流程:
text
用户(Streamlit 前端)
│
▼
FastAPI 异步后端
│
├── RAG 模块(Chroma + BGE Embedding)
│ └── 文档检索(故障手册、操作 SOP)
│
├── Agent 模块(LangGraph)
│ ├── Planner:任务拆解
│ ├── Executor:调用工具(Tools)
│ └── Reflector:结果反思与重试
│
├── 工具集(Tools)
│ ├── WebSearch(模拟)
│ ├── SQL_Query(数据库查询)
│ ├── Execute_Script(Shell 执行)
│ └── Send_Email(邮件通知)
│
├── 记忆模块(Redis 短期 + Chroma 长期)
│
└── 模型层(支持 OpenAI / Ollama / vLLM)技术栈:
项目结构:
ai_ops_agent/
├── backend/
│ ├── app/
│ │ ├── __init__.py
│ │ ├── main.py # FastAPI 入口
│ │ ├── config.py # 配置管理
│ │ ├── rag_engine.py # RAG 检索与索引
│ │ ├── agent_graph.py # LangGraph Agent 定义
│ │ ├── tools.py # 工具函数实现
│ │ ├── memory.py # 短期记忆管理
│ │ └── models.py # Pydantic 模型
│ ├── data/ # 知识库文档(PDF/TXT/MD)
│ ├── chroma_db/ # 向量库持久化目录
│ ├── requirements.txt
│ └── Dockerfile
├── frontend/
│ ├── streamlit_app.py
│ ├── requirements.txt
│ └── Dockerfile
├── docker-compose.yml
└── .env.examplebackend/requirements.txt:
fastapi==0.115.0
uvicorn[standard]==0.30.0
langchain==0.3.0
langchain-community==0.3.0
langgraph==0.2.0
chromadb==0.5.0
sentence-transformers==3.0.0
openai==1.40.0
pypdf==5.0.0
python-multipart==0.0.9
redis==5.0.0
python-dotenv==1.0.0
httpx==0.27.0frontend/requirements.txt:
streamlit==1.37.0
requests==2.32.0config.py)import os
from dotenv import load_dotenv
load_dotenv()
class Settings:
# LLM 配置
LLM_PROVIDER = os.getenv("LLM_PROVIDER", "openai") # openai | ollama | vllm
OPENAI_API_KEY = os.getenv("OPENAI_API_KEY", "")
OPENAI_BASE_URL = os.getenv("OPENAI_BASE_URL", "https://api.openai.com/v1")
LLM_MODEL = os.getenv("LLM_MODEL", "gpt-3.5-turbo")
# Ollama(本地)
OLLAMA_BASE_URL = os.getenv("OLLAMA_BASE_URL", "http://localhost:11434")
OLLAMA_MODEL = os.getenv("OLLAMA_MODEL", "qwen2.5:7b")
# Embedding
EMBEDDING_MODEL = os.getenv("EMBEDDING_MODEL", "BAAI/bge-small-zh-v1.5")
CHROMA_PERSIST_DIR = os.getenv("CHROMA_PERSIST_DIR", "./chroma_db")
# RAG 参数
TOP_K = int(os.getenv("TOP_K", "5"))
CHUNK_SIZE = int(os.getenv("CHUNK_SIZE", "500"))
CHUNK_OVERLAP = int(os.getenv("CHUNK_OVERLAP", "50"))
# Redis(短期记忆)
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
# 工具配置
EXEC_SCRIPT_ALLOWED = os.getenv("EXEC_SCRIPT_ALLOWED", "false").lower() == "true"
SMTP_HOST = os.getenv("SMTP_HOST", "smtp.example.com")
SMTP_PORT = int(os.getenv("SMTP_PORT", "587"))
SMTP_USER = os.getenv("SMTP_USER", "")
SMTP_PASS = os.getenv("SMTP_PASS", "")
settings = Settings()rag_engine.py)支持文档上传、切分、向量化、检索。
from langchain_community.document_loaders import PyPDFLoader, TextLoader, UnstructuredMarkdownLoader
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_community.embeddings import HuggingFaceEmbeddings
from langchain_community.vectorstores import Chroma
from typing import List, Optional
import os
from .config import settings
class RAGEngine:
def __init__(self):
self.embeddings = HuggingFaceEmbeddings(
model_name=settings.EMBEDDING_MODEL,
model_kwargs={'device': 'cpu'}, # 可改 cuda
encode_kwargs={'normalize_embeddings': True}
)
self.vectorstore = None
self.retriever = None
self._load_existing()
def _load_existing(self):
if os.path.exists(settings.CHROMA_PERSIST_DIR) and os.listdir(settings.CHROMA_PERSIST_DIR):
self.vectorstore = Chroma(
persist_directory=settings.CHROMA_PERSIST_DIR,
embedding_function=self.embeddings
)
self.retriever = self.vectorstore.as_retriever(
search_type="similarity",
search_kwargs={"k": settings.TOP_K}
)
return True
return False
def load_and_split_documents(self, file_paths: List[str]):
"""加载并切分文档"""
all_docs = []
for path in file_paths:
if path.endswith('.pdf'):
loader = PyPDFLoader(path)
elif path.endswith('.md'):
loader = UnstructuredMarkdownLoader(path)
else:
loader = TextLoader(path, encoding='utf-8')
all_docs.extend(loader.load())
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=settings.CHUNK_SIZE,
chunk_overlap=settings.CHUNK_OVERLAP,
separators=["\n\n", "\n", "。", "!", "?", ";", ",", " ", ""]
)
return text_splitter.split_documents(all_docs)
def rebuild_index(self, docs):
"""重建向量索引"""
self.vectorstore = Chroma.from_documents(
documents=docs,
embedding=self.embeddings,
persist_directory=settings.CHROMA_PERSIST_DIR
)
self.vectorstore.persist()
self.retriever = self.vectorstore.as_retriever(
search_type="similarity",
search_kwargs={"k": settings.TOP_K}
)
def retrieve(self, query: str) -> List[str]:
"""检索相关文档内容"""
if not self.retriever:
return []
docs = self.retriever.invoke(query)
return [doc.page_content for doc in docs]
def retrieve_with_sources(self, query: str):
if not self.retriever:
return []
docs = self.retriever.invoke(query)
return [{"content": doc.page_content, "source": doc.metadata.get("source", "unknown")} for doc in docs]
# 全局单例
rag_engine = RAGEngine()tools.py)Agent 可调用的工具集。
from langchain_core.tools import tool
import subprocess
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import httpx
import json
from .config import settings
@tool
def web_search(query: str) -> str:
"""模拟网络搜索,用于获取实时信息(如天气、新闻)"""
# 实际可接入 SerpAPI、Bing Search 等
return f"搜索结果:关于 '{query}' 的最新信息,建议查阅官方文档。"
@tool
def query_database(sql: str) -> str:
"""执行 SQL 查询(只读),用于获取业务数据,如订单状态、服务器指标"""
# 实际接入真实数据库,此处模拟
if "select" in sql.lower():
return f"模拟查询结果:{sql} 返回 3 条记录。"
return "仅支持 SELECT 查询"
@tool
def execute_shell_command(command: str) -> str:
"""执行 Shell 命令(谨慎使用,仅限授权环境)"""
if not settings.EXEC_SCRIPT_ALLOWED:
return "执行脚本功能未启用,请在配置中开启。"
try:
result = subprocess.run(command, shell=True, capture_output=True, text=True, timeout=10)
if result.returncode == 0:
return f"命令执行成功:\n{result.stdout}"
else:
return f"命令执行失败:\n{result.stderr}"
except Exception as e:
return f"执行异常:{str(e)}"
@tool
def send_email(recipient: str, subject: str, body: str) -> str:
"""发送邮件通知"""
if not settings.SMTP_USER:
return "邮件服务未配置"
try:
msg = MIMEMultipart()
msg['From'] = settings.SMTP_USER
msg['To'] = recipient
msg['Subject'] = subject
msg.attach(MIMEText(body, 'plain'))
with smtplib.SMTP(settings.SMTP_HOST, settings.SMTP_PORT) as server:
server.starttls()
server.login(settings.SMTP_USER, settings.SMTP_PASS)
server.send_message(msg)
return f"邮件已发送至 {recipient}"
except Exception as e:
return f"邮件发送失败:{str(e)}"
@tool
def check_service_status(service_name: str) -> str:
"""检查系统服务运行状态(通过模拟接口)"""
# 实际可调用 Prometheus API 或 systemctl
statuses = {"order-service": "running", "payment-service": "degraded", "inventory-service": "stopped"}
status = statuses.get(service_name, "unknown")
return f"服务 {service_name} 状态:{status}"
tools_list = [web_search, query_database, execute_shell_command, send_email, check_service_status]agent_graph.py)使用 LangGraph 构建 ReAct 风格的 Agent,支持循环推理和工具调用。
from typing import Literal, List, Dict, Any
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolExecutor
from langchain_core.messages import HumanMessage, AIMessage, SystemMessage, ToolMessage
from langchain_openai import ChatOpenAI
from langchain_ollama import ChatOllama
import json
from .config import settings
from .tools import tools_list
from .rag_engine import rag_engine
# 初始化 LLM
def init_llm():
if settings.LLM_PROVIDER == "openai":
return ChatOpenAI(
api_key=settings.OPENAI_API_KEY,
base_url=settings.OPENAI_BASE_URL,
model=settings.LLM_MODEL,
temperature=0.1,
streaming=True
)
elif settings.LLM_PROVIDER == "ollama":
return ChatOllama(
base_url=settings.OLLAMA_BASE_URL,
model=settings.OLLAMA_MODEL,
temperature=0.1,
)
else:
raise ValueError("Unsupported LLM provider")
llm = init_llm()
tool_executor = ToolExecutor(tools_list)
# 绑定工具到 LLM
llm_with_tools = llm.bind_tools(tools_list)
# 定义状态
class AgentState(Dict[str, Any]):
messages: List[Any]
current_task: str
iteration: int
# 节点函数
def call_model(state: AgentState):
messages = state["messages"]
# 如果当前问题是知识类,先检索 RAG 增强上下文
if "故障" in state.get("current_task", "") or "手册" in state.get("current_task", ""):
query = state["current_task"]
docs = rag_engine.retrieve(query)
if docs:
context = "\n\n".join(docs[:2])
system_msg = SystemMessage(content=f"参考以下运维知识库内容回答:\n{context}")
messages = [system_msg] + messages
response = llm_with_tools.invoke(messages)
return {"messages": messages + [response]}
def call_tool(state: AgentState):
last_message = state["messages"][-1]
tool_calls = last_message.tool_calls
results = []
for tool_call in tool_calls:
tool_name = tool_call["name"]
tool_args = tool_call["args"]
# 执行工具
result = tool_executor.invoke(tool_call)
results.append(ToolMessage(
content=result,
tool_call_id=tool_call["id"]
))
return {"messages": state["messages"] + results}
def should_continue(state: AgentState) -> Literal["tools", "__end__"]:
last_message = state["messages"][-1]
if hasattr(last_message, "tool_calls") and last_message.tool_calls:
return "tools"
return "__end__"
# 构建图
workflow = StateGraph(AgentState)
workflow.add_node("agent", call_model)
workflow.add_node("tools", call_tool)
workflow.set_entry_point("agent")
workflow.add_conditional_edges(
"agent",
should_continue,
{
"tools": "tools",
"__end__": "__end__"
}
)
workflow.add_edge("tools", "agent")
agent_graph = workflow.compile()
# 执行函数(供 API 调用)
def run_agent(question: str, session_id: str = None) -> str:
"""同步执行 Agent 并返回最终答案"""
state = {
"messages": [HumanMessage(content=question)],
"current_task": question,
"iteration": 0
}
# 最多迭代 5 轮
for _ in range(5):
result = agent_graph.invoke(state)
state = result
last_msg = state["messages"][-1]
if not (hasattr(last_msg, "tool_calls") and last_msg.tool_calls):
break
# 提取最终回答
for msg in reversed(state["messages"]):
if isinstance(msg, AIMessage) and msg.content:
return msg.content
return "无法生成回答"
async def run_agent_stream(question: str, session_id: str = None):
"""流式执行 Agent,逐 token 返回"""
# 简化:直接调用 LLM 流式(实际可按需支持工具中途输出)
# 为演示,我们先用普通模式返回完整结果,再逐字流式
result = run_agent(question, session_id)
for char in result:
yield char
# 模拟异步等待
import asyncio
await asyncio.sleep(0.01)注:LangGraph 的流式支持更精细,此处简化以便演示完整流程。
memory.py)基于 Redis 存储会话上下文。
import redis
import json
from .config import settings
redis_client = redis.Redis.from_url(settings.REDIS_URL, decode_responses=True)
def get_session_history(session_id: str, limit: int = 10):
"""获取会话最近 N 条消息"""
key = f"session:{session_id}"
messages = redis_client.lrange(key, -limit, -1)
return [json.loads(m) for m in messages]
def add_message(session_id: str, role: str, content: str):
key = f"session:{session_id}"
redis_client.rpush(key, json.dumps({"role": role, "content": content}))
# 保留最近 100 条
redis_client.ltrim(key, -100, -1)
def clear_session(session_id: str):
redis_client.delete(f"session:{session_id}")main.py)提供文档上传、对话、会话管理等 API。
from fastapi import FastAPI, UploadFile, File, HTTPException, Depends
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
from typing import Optional
import shutil
import tempfile
import os
from .rag_engine import rag_engine
from .agent_graph import run_agent, run_agent_stream
from .memory import get_session_history, add_message, clear_session
from .config import settings
app = FastAPI(title="AI 运维助手 (RAG + Agent)", version="1.0")
class ChatRequest(BaseModel):
question: str
session_id: Optional[str] = "default"
@app.on_event("startup")
async def startup():
# 确保向量库目录存在
os.makedirs(settings.CHROMA_PERSIST_DIR, exist_ok=True)
if rag_engine.vectorstore:
print("✅ RAG 引擎加载成功")
else:
print("⚠️ 未找到向量库,请上传文档")
@app.post("/upload")
async def upload_knowledge(files: list[UploadFile] = File(...)):
"""上传知识库文档(PDF/TXT/MD)"""
temp_dir = tempfile.mkdtemp()
saved_paths = []
for f in files:
if not f.filename.endswith(('.pdf', '.txt', '.md')):
raise HTTPException(400, f"不支持类型: {f.filename}")
path = os.path.join(temp_dir, f.filename)
with open(path, "wb") as buffer:
shutil.copyfileobj(f.file, buffer)
saved_paths.append(path)
try:
docs = rag_engine.load_and_split_documents(saved_paths)
rag_engine.rebuild_index(docs)
return {"message": f"成功处理 {len(saved_paths)} 个文档,共 {len(docs)} 个切片"}
except Exception as e:
raise HTTPException(500, f"构建向量库失败: {str(e)}")
finally:
shutil.rmtree(temp_dir, ignore_errors=True)
@app.post("/chat")
async def chat(request: ChatRequest):
"""Agent 对话接口(流式)"""
# 记录用户问题(短期记忆)
add_message(request.session_id, "user", request.question)
# 获取历史记忆(可注入到 Agent)
history = get_session_history(request.session_id)
# 注:完整实现可将 history 构建为消息列表传递给 Agent
def generate():
full_answer = ""
try:
async for chunk in run_agent_stream(request.question, request.session_id):
full_answer += chunk
yield chunk
# 保存助手回复
add_message(request.session_id, "assistant", full_answer)
except Exception as e:
yield f"错误: {str(e)}"
return StreamingResponse(generate(), media_type="text/plain; charset=utf-8")
@app.get("/history/{session_id}")
async def get_history(session_id: str):
"""获取会话历史"""
return get_session_history(session_id)
@app.delete("/history/{session_id}")
async def clear_history(session_id: str):
clear_session(session_id)
return {"message": "会话已清空"}
@app.get("/health")
async def health():
return {"status": "ok", "rag_ready": rag_engine.vectorstore is not None}cd backend
uvicorn app.main:app --host 0.0.0.0 --port 8000 --reload提供简洁的聊天界面,支持文档上传和会话管理。
# frontend/streamlit_app.py
import streamlit as st
import requests
import json
API_BASE = "http://localhost:8000"
st.set_page_config(page_title="AI 运维助手", page_icon="🤖")
st.title("🤖 企业智能运维助手 (RAG + Agent)")
# 初始化 session
if "messages" not in st.session_state:
st.session_state.messages = []
if "session_id" not in st.session_state:
st.session_state.session_id = "default"
# 侧边栏
with st.sidebar:
st.header("知识库管理")
uploaded_files = st.file_uploader("上传运维手册/故障记录 (PDF/TXT/MD)",
accept_multiple_files=True,
type=["pdf", "txt", "md"])
if uploaded_files and st.button("更新知识库"):
files = [("files", (f.name, f, f.type)) for f in uploaded_files]
with st.spinner("处理文档..."):
resp = requests.post(f"{API_BASE}/upload", files=files)
if resp.status_code == 200:
st.success(resp.json()["message"])
else:
st.error(f"失败: {resp.text}")
st.divider()
if st.button("清空会话"):
requests.delete(f"{API_BASE}/history/{st.session_state.session_id}")
st.session_state.messages = []
st.rerun()
st.caption("会话 ID: " + st.session_state.session_id)
# 显示历史消息
for msg in st.session_state.messages:
with st.chat_message(msg["role"]):
st.markdown(msg["content"])
# 输入
user_input = st.chat_input("描述故障或提问...")
if user_input:
# 显示用户消息
st.session_state.messages.append({"role": "user", "content": user_input})
with st.chat_message("user"):
st.markdown(user_input)
# 调用流式 API
with st.chat_message("assistant"):
placeholder = st.empty()
full_response = ""
try:
resp = requests.post(
f"{API_BASE}/chat",
json={"question": user_input, "session_id": st.session_state.session_id},
stream=True
)
for chunk in resp.iter_content(chunk_size=128, decode_unicode=True):
if chunk:
full_response += chunk
placeholder.markdown(full_response + "▌")
placeholder.markdown(full_response)
except Exception as e:
placeholder.error(f"请求异常: {e}")
st.session_state.messages.append({"role": "assistant", "content": full_response})运行前端:
streamlit run frontend/streamlit_app.py为了生产环境可用,需增强 Agent 的可观测性(中间步骤展示)。我们可在 run_agent_stream 中 yield 结构化数据(如工具调用日志),前端解析展示。
改进版流式生成器(片段):
async def run_agent_stream_verbose(question: str, session_id: str):
# 模拟步骤
yield json.dumps({"type": "status", "msg": "正在分析问题..."})
# 调用 RAG
docs = rag_engine.retrieve(question)
if docs:
yield json.dumps({"type": "rag", "docs": docs[:2]})
# 调用 Agent
result = run_agent(question, session_id)
yield json.dumps({"type": "answer", "content": result})前端相应解析并显示步骤折叠面板,提升用户体验。
docker-compose.yml:
version: '3.8'
services:
redis:
image: redis:7-alpine
ports:
- "6379:6379"
restart: unless-stopped
backend:
build: ./backend
ports:
- "8000:8000"
environment:
- REDIS_URL=redis://redis:6379/0
- OPENAI_API_KEY=${OPENAI_API_KEY}
- LLM_PROVIDER=${LLM_PROVIDER:-openai}
volumes:
- ./backend/chroma_db:/app/chroma_db
- ./backend/data:/app/data
depends_on:
- redis
restart: unless-stopped
frontend:
build: ./frontend
ports:
- "8501:8501"
depends_on:
- backend
restart: unless-stoppedbackend/Dockerfile:
FROM python:3.10-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]frontend/Dockerfile:
FROM python:3.10-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
CMD ["streamlit", "run", "streamlit_app.py", "--server.port=8501", "--server.address=0.0.0.0"]启动:
docker-compose up -d本文完整构建了一个企业级 RAG + Agent 混合智能运维助手,核心亮点:
可扩展方向:
AI 全栈工程师的核心价值,在于将模型智能与工程能力结合,交付真正解决业务问题的系统。本文所有代码均可运行,建议根据自身场景替换实际工具和知识库,快速落地。
免责声明:本文代码仅供学习参考,生产环境需结合权限管理、日志审计、安全沙箱等措施。
关于作者:AI 应用架构师,专注于大模型工程化与智能体系统设计。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。