从 0 到 1 构建企业级 AI Agent:FastAPI + RAG + Milvus + MCP 全流程实战
本文通过一个可运行的企业知识库助手,串联 FastAPI、RAG、Milvus、Function Calling 和 MCP。读完后,你将得到一个可以回答企业内部文档问题、查询业务工具并通过 HTTP 接口对外提供服务的 AI Agent。
一、项目效果
我们要构建的系统包含三个能力:
- 从企业文档中检索答案,降低大模型幻觉。
- 通过 MCP 调用订单查询等外部业务工具。
- 使用 FastAPI 提供统一的对话接口。
整体调用链路如下:
二、技术选型
| 模块 | 技术 | 作用 |
|---|---|---|
| Web 服务 | FastAPI | 提供聊天接口和健康检查接口 |
| 大模型 | OpenAI 兼容 API | 负责理解问题和选择工具 |
| 知识库 | RAG | 检索企业文档并注入上下文 |
| 向量数据库 | Milvus / Milvus Lite | 保存和检索文档向量 |
| 工具协议 | MCP | 标准化暴露外部工具 |
| 文档解析 | pypdf | 读取 PDF 文档 |
示例代码默认使用 OpenAI 兼容接口。只要服务商支持 Chat Completions 和 Embeddings 接口,就可以通过环境变量替换地址和模型。
三、项目结构
ai-agent-demo/
├── app/
│ ├── __init__.py
│ ├── agent.py
│ ├── main.py
│ ├── mcp_bridge.py
│ └── rag.py
├── data/
│ └── employee-handbook.pdf
├── .env.example
├── ingest.py
├── mcp_server.py
└── requirements.txt
四、安装依赖
创建 requirements.txt:
fastapi
uvicorn[standard]
python-dotenv
openai
pymilvus
pypdf
mcp
安装依赖:
python -m venv .venv
# Windows
.venv\Scripts\activate
# macOS/Linux
source .venv/bin/activate
pip install -r requirements.txt
五、配置大模型和 Milvus
创建 .env 文件:
LLM_API_KEY=替换为你的模型服务密钥
LLM_BASE_URL=https://api.openai.com/v1
LLM_MODEL=gpt-4o-mini
EMBEDDING_MODEL=text-embedding-3-small
# 开发环境可以直接使用 Milvus Lite 文件模式
MILVUS_URI=./milvus_demo.db
MILVUS_COLLECTION=company_knowledge
# 使用独立 Milvus 服务时,可以改成:
# MILVUS_URI=http://127.0.0.1:19530
Milvus Lite 适合本地开发和文章演示;生产环境建议使用独立 Milvus 集群,并把 MILVUS_URI 改成服务地址。
六、实现 RAG 知识库
RAG 的核心流程是:
文档 -> 文本切分 -> Embedding -> Milvus
用户问题 -> Embedding -> 相似度检索 -> 大模型生成答案
创建 app/rag.py:
import os
from typing import Iterable
from dotenv import load_dotenv
from openai import OpenAI
from pymilvus import DataType, MilvusClient
load_dotenv()
class RAGStore:
def __init__(self) -> None:
self.collection = os.getenv("MILVUS_COLLECTION", "company_knowledge")
self.embedding_model = os.getenv(
"EMBEDDING_MODEL", "text-embedding-3-small"
)
self.embedding_client = OpenAI(
api_key=os.environ["LLM_API_KEY"],
base_url=os.getenv("LLM_BASE_URL"),
)
self.client = MilvusClient(
uri=os.getenv("MILVUS_URI", "./milvus_demo.db")
)
self._ensure_collection()
def _ensure_collection(self) -> None:
if self.client.has_collection(self.collection):
return
schema = self.client.create_schema(
auto_id=True,
enable_dynamic_field=False,
)
schema.add_field(
field_name="id",
datatype=DataType.INT64,
is_primary=True,
)
schema.add_field(
field_name="vector",
datatype=DataType.FLOAT_VECTOR,
dim=1536,
)
schema.add_field(
field_name="text",
datatype=DataType.VARCHAR,
max_length=65535,
)
schema.add_field(
field_name="source",
datatype=DataType.VARCHAR,
max_length=1024,
)
index_params = self.client.prepare_index_params()
index_params.add_index(
field_name="vector",
index_type="AUTOINDEX",
metric_type="COSINE",
)
self.client.create_collection(
collection_name=self.collection,
schema=schema,
index_params=index_params,
)
def embed(self, texts: list[str]) -> list[list[float]]:
response = self.embedding_client.embeddings.create(
model=self.embedding_model,
input=texts,
)
return [item.embedding for item in response.data]
def upsert(self, chunks: Iterable[tuple[str, str]]) -> int:
records = list(chunks)
if not records:
return 0
texts = [item[0] for item in records]
sources = [item[1] for item in records]
vectors = self.embed(texts)
rows = [
{"vector": vector, "text": text, "source": source}
for vector, text, source in zip(vectors, texts, sources)
]
self.client.insert(collection_name=self.collection, data=rows)
return len(rows)
def search(self, question: str, top_k: int = 4) -> str:
query_vector = self.embed([question])[0]
result = self.client.search(
collection_name=self.collection,
data=[query_vector],
anns_field="vector",
limit=top_k,
output_fields=["text", "source"],
search_params={"metric_type": "COSINE", "params": {}},
)
hits = []
for hit in result[0]:
entity = hit["entity"]
hits.append(
{
"score": round(float(hit["distance"]), 4),
"source": entity.get("source", "unknown"),
"text": entity.get("text", ""),
}
)
return str(hits)
rag_store = RAGStore()
这里将向量维度写成了 1536,因为示例使用的是 text-embedding-3-small。如果更换 Embedding 模型,必须同步修改 Milvus 字段的 dim,否则插入向量时会报维度不匹配。
七、准备文档并写入知识库
创建 ingest.py:
from pathlib import Path
from pypdf import PdfReader
from app.rag import rag_store
def split_text(text: str, chunk_size: int = 800, overlap: int = 120):
text = " ".join(text.split())
start = 0
while start < len(text):
end = min(start + chunk_size, len(text))
chunk = text[start:end].strip()
if chunk:
yield chunk
if end == len(text):
break
start = end - overlap
def read_pdf(path: Path):
reader = PdfReader(str(path))
return "\n".join(page.extract_text() or "" for page in reader.pages)
def main():
records = []
for path in Path("data").glob("*.pdf"):
content = read_pdf(path)
records.extend((chunk, path.name) for chunk in split_text(content))
count = rag_store.upsert(records)
print(f"已写入 {count} 个知识片段")
if __name__ == "__main__":
main()
执行导入:
python ingest.py
生产环境建议进一步增加标题感知切分、表格解析、重复文档过滤和增量更新机制。固定字符数切分适合验证链路,不一定适合所有业务文档。
八、使用 MCP 暴露业务工具
RAG 适合回答“文档里写了什么”,但不适合回答“订单现在是什么状态”。后者应该调用真实业务系统。MCP 可以把这类能力以标准工具形式暴露给 Agent。
创建 mcp_server.py:
from mcp.server.fastmcp import FastMCP
mcp = FastMCP("company-business-tools")
@mcp.tool()
def get_order_status(order_id: str) -> dict:
"""查询订单状态。实际项目中应替换为数据库或业务 API 调用。"""
demo_orders = {
"A10001": {"status": "已发货", "物流": "SF123456789"},
"A10002": {"status": "待支付", "物流": ""},
}
return demo_orders.get(
order_id,
{"status": "未找到订单", "物流": ""},
)
@mcp.tool()
def get_service_notice() -> str:
"""获取当前服务公告。"""
return "系统服务时间为工作日 9:00-18:00。"
if __name__ == "__main__":
mcp.run(transport="stdio")
注意:示例中的订单数据只是演示。真实项目中,MCP 工具必须增加身份校验、权限控制、参数校验和操作审计,不能让模型直接获得无限制的数据库权限。
九、让 Agent 动态加载 MCP 工具
创建 app/mcp_bridge.py:
import os
import sys
from contextlib import AsyncExitStack
from pathlib import Path
from mcp import ClientSession, StdioServerParameters
from mcp.client.stdio import stdio_client
class MCPToolBridge:
def __init__(self) -> None:
self.stack = AsyncExitStack()
self.session = None
self.openai_tools = []
async def connect(self) -> None:
server_path = Path(__file__).resolve().parents[1] / "mcp_server.py"
params = StdioServerParameters(
command=sys.executable,
args=[str(server_path)],
env=os.environ.copy(),
)
read_stream, write_stream = await self.stack.enter_async_context(
stdio_client(params)
)
self.session = await self.stack.enter_async_context(
ClientSession(read_stream, write_stream)
)
await self.session.initialize()
response = await self.session.list_tools()
self.openai_tools = [
{
"type": "function",
"function": {
"name": tool.name,
"description": tool.description or "",
"parameters": tool.inputSchema
or {"type": "object", "properties": {}},
},
}
for tool in response.tools
]
async def call(self, name: str, arguments: dict) -> str:
result = await self.session.call_tool(name, arguments)
parts = []
for item in result.content:
if hasattr(item, "text"):
parts.append(item.text)
else:
parts.append(str(item))
return "\n".join(parts)
async def close(self) -> None:
await self.stack.aclose()
这段桥接代码做了三件事:启动 MCP Server、读取工具描述、把 MCP 工具转换成大模型可以理解的工具定义。这样新增 MCP 工具时,Agent 不需要手工修改工具列表。
十、实现 Agent 工具调用循环
创建 app/agent.py:
import json
import os
from dotenv import load_dotenv
from openai import AsyncOpenAI
from app.mcp_bridge import MCPToolBridge
from app.rag import rag_store
load_dotenv()
llm = AsyncOpenAI(
api_key=os.environ["LLM_API_KEY"],
base_url=os.getenv("LLM_BASE_URL"),
)
LOCAL_TOOLS = [
{
"type": "function",
"function": {
"name": "search_knowledge_base",
"description": "检索企业内部文档,适合回答制度、产品和流程问题。",
"parameters": {
"type": "object",
"properties": {
"question": {
"type": "string",
"description": "需要检索的问题",
},
"top_k": {
"type": "integer",
"description": "返回的知识片段数量,默认 4",
},
},
"required": ["question"],
},
},
}
]
SYSTEM_PROMPT = """
你是企业智能助手。请遵守以下规则:
1. 涉及企业内部知识时,必须先调用 search_knowledge_base。
2. 涉及订单、服务状态等实时信息时,必须调用对应 MCP 工具。
3. 不要编造工具没有返回的事实。
4. 如果检索结果不足以回答,请明确说明信息不足。
5. 回答简洁,并在知识库答案后标注来源文件名。
"""
async def run_agent(question: str, bridge: MCPToolBridge) -> str:
messages = [
{"role": "system", "content": SYSTEM_PROMPT},
{"role": "user", "content": question},
]
tools = LOCAL_TOOLS + bridge.openai_tools
for _ in range(5):
response = await llm.chat.completions.create(
model=os.getenv("LLM_MODEL", "gpt-4o-mini"),
messages=messages,
tools=tools,
tool_choice="auto",
temperature=0.1,
)
message = response.choices[0].message
messages.append(message.model_dump(exclude_none=True))
if not message.tool_calls:
return message.content or "暂时无法生成答案。"
for tool_call in message.tool_calls:
name = tool_call.function.name
arguments = json.loads(tool_call.function.arguments or "{}")
if name == "search_knowledge_base":
result = rag_store.search(
question=arguments["question"],
top_k=arguments.get("top_k", 4),
)
else:
result = await bridge.call(name, arguments)
messages.append(
{
"role": "tool",
"tool_call_id": tool_call.id,
"content": result,
}
)
return "工具调用次数超过限制,请稍后重试。"
Agent 的关键不是简单地调用一次大模型,而是实现一个有限循环:
用户问题 -> 大模型判断 -> 调用工具 -> 返回工具结果 -> 大模型总结
代码中的最大循环次数是 5,可以防止模型因为工具选择异常而无限调用。生产环境还应增加单次请求 Token 限制、超时控制和费用统计。
十一、使用 FastAPI 对外提供服务
创建 app/main.py:
from contextlib import asynccontextmanager
from fastapi import FastAPI
from pydantic import BaseModel, Field
from app.agent import run_agent
from app.mcp_bridge import MCPToolBridge
bridge = MCPToolBridge()
@asynccontextmanager
async def lifespan(app: FastAPI):
await bridge.connect()
yield
await bridge.close()
app = FastAPI(
title="Enterprise AI Agent",
version="1.0.0",
lifespan=lifespan,
)
class ChatRequest(BaseModel):
question: str = Field(min_length=1, max_length=2000)
class ChatResponse(BaseModel):
answer: str
@app.get("/health")
async def health():
return {"status": "ok"}
@app.post("/chat", response_model=ChatResponse)
async def chat(request: ChatRequest):
answer = await run_agent(request.question, bridge)
return ChatResponse(answer=answer)
启动服务:
uvicorn app.main:app --reload --port 8000
调用接口:
curl -X POST http://127.0.0.1:8000/chat \
-H "Content-Type: application/json" \
-d '{"question":"订单 A10001 当前是什么状态?"}'
也可以访问 http://127.0.0.1:8000/docs,使用 FastAPI 自动生成的页面测试接口。
十二、一次完整请求是如何执行的
当用户输入“订单 A10001 当前是什么状态”时:
- FastAPI 接收问题并校验参数。
- Agent 判断这是实时业务问题。
- 大模型选择 MCP 工具
get_order_status。 - MCP Server 返回订单状态。
- Agent 根据工具结果生成自然语言答案。
当用户输入“公司的报销流程是什么”时:
- Agent 选择
search_knowledge_base。 - 问题被转换成向量。
- Milvus 返回相似文档片段。
- Agent 根据检索结果回答,并标注来源。
十三、企业级改造重点
1. 优化文档切分
固定字符切分只是起点。企业文档建议按标题、章节、表格和段落切分,并为每个片段保存部门、文档版本和权限标签。
2. 增加混合检索
只使用向量检索可能无法处理订单号、产品编号等精确关键词。可以将 BM25 关键词检索和向量检索结合,再使用重排序模型提升结果质量。
3. 增加权限过滤
检索条件不能只有问题向量,还应包含当前用户的部门、角色和数据权限。未经授权的文档不能因为语义相似而被召回。
4. 控制 MCP 工具风险
查询类工具可以自动调用;删除、退款、审批等写操作应该增加人工确认、二次认证和审计日志。
5. 建立评估集
至少准备以下指标:
- 检索召回率:正确文档是否被找到。
- 答案准确率:回答是否符合事实。
- 引用完整性:答案是否能够追溯到来源。
- 工具选择准确率:Agent 是否调用了正确工具。
- 平均响应时间和单次请求成本。
十四、常见问题
Milvus 报向量维度不匹配怎么办?
检查 Embedding 模型输出维度,并与集合字段的 dim 保持一致。更换模型后,通常需要创建新的集合,不能直接复用旧集合。
Agent 总是调用错误的工具怎么办?
优化工具名称和描述,明确工具适用场景;同时在系统提示词中写清楚“什么问题必须调用什么工具”。对于高风险操作,不要只依赖提示词,应在服务端增加权限和参数校验。
为什么 RAG 仍然会产生幻觉?
常见原因包括文档切分不合理、召回片段不相关、上下文过长和提示词约束不足。可以通过重排序、问题改写、答案引用和评估集逐项定位。
MCP 和普通函数调用有什么区别?
普通函数调用通常需要在每个应用中手工注册工具;MCP 提供了统一的工具发现和调用协议,使工具可以被多个 AI 应用复用。两者可以结合使用:本地简单能力用函数调用,跨系统业务能力用 MCP。
十五、总结
本文完成了一个最小但完整的企业级 AI Agent:
- FastAPI 负责服务化。
- RAG 负责企业知识检索。
- Milvus 负责向量存储。
- MCP 负责连接外部业务工具。
- Function Calling 负责让 Agent 自主选择工具。
真正的企业级系统还需要继续补充权限、监控、评估、缓存、限流、审计和人工确认机制。建议按照“先跑通链路,再逐步增强可靠性”的顺序迭代,而不是一开始就堆叠大量框架。
更多推荐

所有评论(0)