本文通过一个可运行的企业知识库助手,串联 FastAPI、RAG、Milvus、Function Calling 和 MCP。读完后,你将得到一个可以回答企业内部文档问题、查询业务工具并通过 HTTP 接口对外提供服务的 AI Agent。

一、项目效果

我们要构建的系统包含三个能力:

  1. 从企业文档中检索答案,降低大模型幻觉。
  2. 通过 MCP 调用订单查询等外部业务工具。
  3. 使用 FastAPI 提供统一的对话接口。

整体调用链路如下:

知识问题

业务操作

用户问题

FastAPI

AI Agent

RAG 检索

Milvus 向量库

MCP Server

业务系统

大模型

二、技术选型

模块 技术 作用
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 当前是什么状态”时:

  1. FastAPI 接收问题并校验参数。
  2. Agent 判断这是实时业务问题。
  3. 大模型选择 MCP 工具 get_order_status
  4. MCP Server 返回订单状态。
  5. Agent 根据工具结果生成自然语言答案。

当用户输入“公司的报销流程是什么”时:

  1. Agent 选择 search_knowledge_base
  2. 问题被转换成向量。
  3. Milvus 返回相似文档片段。
  4. 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 自主选择工具。

真正的企业级系统还需要继续补充权限、监控、评估、缓存、限流、审计和人工确认机制。建议按照“先跑通链路,再逐步增强可靠性”的顺序迭代,而不是一开始就堆叠大量框架。

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐