这次我们来看一个结合 LangGraph 和 MCP 的大模型实战开发项目。这个教程来自马士兵教育的 AI 大模型课程,重点不是讲抽象概念,而是手把手教你搭建一个能实际运行的智能小秘书系统。如果你正在学习 Agent 开发、多工具链集成,或者想了解如何用 LangGraph 编排复杂工作流,这篇文章会直接带你看环境准备、代码实现和效果验证。
LangGraph 是 LangChain 团队推出的工作流编排框架,专门解决多步骤、有状态、带循环的 Agent 任务。MCP(Model Context Protocol)则是一种新的工具调用协议,让大模型能更安全、标准化地使用外部工具。两者结合,可以构建出既能理解复杂指令,又能调用各种 API 和本地工具的智能助手。
本文会重点演示几个核心环节:环境怎么配、Graph 怎么画、MCP Server 怎么接、任务状态怎么流转,以及最终怎么让这个“小秘书”帮你查天气、读文档、订日程。我们会用可运行的代码示例说明每个环节,并给出本地测试时常见的端口冲突、依赖版本问题排查方法。
1. 核心能力速览
| 能力项 | 说明 |
|---|---|
| 技术栈 | LangGraph(工作流编排) + MCP(工具协议) + 大模型(如 OpenAI GPT-4、Claude 或本地模型) |
| 主要功能 | 多步骤任务规划、工具调用、状态持久化、循环控制、错误处理 |
| 硬件门槛 | 依赖所选大模型;若用本地模型需 GPU,若用 API 则只需 CPU 和网络 |
| 启动方式 | Python 脚本启动、Jupyter Notebook 交互测试、FastAPI 服务化 |
| 接口能力 | 支持 HTTP API 调用,可集成到 Web、移动端或第三方系统 |
| 批量任务 | 通过工作流状态管理支持批量异步处理 |
| 适合场景 | 智能助手、自动化流程、数据查询与分析、多工具协作任务 |
2. 适用场景与使用边界
这个 LangGraph + MCP 的智能小秘书最适合以下几类场景:
适合场景:
- 个人效率助手:帮你汇总日程、查询信息、自动整理文档
- 企业内部助手:连接公司内部的 API(如 CRM、项目管理工具),完成数据查询、报表生成
- 开发测试助手:集成代码检查、日志查询、部署状态跟踪等工具链
- 数据分析助手:调用数据库、可视化工具,完成多步骤分析任务
使用边界:
- 工具调用依赖 MCP Server 的可用性和权限设置,需要提前配置好工具端点
- 复杂工作流可能会涉及多次大模型调用,API 成本或本地推理资源需合理规划
- 任务执行过程中如果外部工具失效,需要有重试或降级方案
- 涉及用户隐私或敏感数据的工具调用,必须做好授权验证和访问日志
3. 环境准备与前置条件
开始前,请确保你的开发环境满足以下条件:
操作系统:Windows 10/11、macOS 或 Linux(推荐 Ubuntu 20.04+)
Python 版本:3.9 或 3.10(避免使用 3.11+,因部分依赖可能尚未完全兼容)
包管理工具:pip 或 conda
核心依赖包:
# 基础框架 langgraph >= 0.0.40 langchain >= 0.1.0 mcp >= 1.0.0 # 大模型接入(根据实际使用的模型选择) openai >= 1.0.0 # 若使用 OpenAI API anthropic >= 0.25.0 # 若使用 Claude # 服务化与交互 fastapi >= 0.104.0 uvicorn >= 0.24.0 jupyter >= 1.0.0 # 用于交互式测试可选本地模型部署: 如果你打算用本地大模型(如 Ollama 部署的 Llama、书生·浦语等),还需要:
- Ollama 或类似本地模型服务
- 足够的 GPU 显存(至少 8GB 用于 7B 模型,16GB+ 用于 13B+ 模型)
4. 安装部署与启动方式
4.1 创建并激活虚拟环境
# 创建虚拟环境 python -m venv langgraph-mcp-demo # 激活环境(Windows) langgraph-mcp-demo\Scripts\activate # 激活环境(macOS/Linux) source langgraph-mcp-demo/bin/activate4.2 安装核心依赖
pip install langgraph langchain mcp openai fastapi uvicorn4.3 配置模型访问凭证
创建.env文件保存你的模型 API 密钥:
# 使用 OpenAI GPT-4 OPENAI_API_KEY=your_openai_api_key_here # 或使用 Anthropic Claude ANTHROPIC_API_KEY=your_anthropic_api_key_here # 若用本地模型,配置基础地址 LOCAL_MODEL_BASE_URL=http://localhost:114344.4 基础服务启动
创建一个最简单的 LangGraph 工作流验证环境:
# demo_verification.py import asyncio from langgraph.graph import StateGraph, END from typing import TypedDict class AgentState(TypedDict): input: str response: str def call_model(state: AgentState): return {"response": f"Received: {state['input']}"} # 构建图 graph_builder = StateGraph(AgentState) graph_builder.add_node("model", call_model) graph_builder.set_entry_point("model") graph_builder.set_finish_point("model") graph = graph_builder.compile() # 测试运行 result = graph.invoke({"input": "Hello LangGraph!"}) print(result["response"])运行测试:
python demo_verification.py如果输出Received: Hello LangGraph!说明基础环境正常。
5. 功能测试与效果验证
5.1 基础工作流测试
我们先构建一个包含多个节点的简单工作流,模拟智能小秘书的核心流程:
# basic_workflow.py from langgraph.graph import StateGraph, END from typing import TypedDict, Annotated from typing_extensions import TypedDict import operator class WorkflowState(TypedDict): user_input: str task_type: Annotated[str, operator.add] tool_calls: Annotated[list, operator.add] final_response: str def classify_task(state: WorkflowState): input_text = state["user_input"].lower() if "天气" in input_text: return {"task_type": "weather_query"} elif "日程" in input_text or "安排" in input_text: return {"task_type": "schedule_manage"} else: return {"task_type": "general_query"} def call_weather_tool(state: WorkflowState): # 模拟天气查询工具调用 return {"tool_calls": ["weather_api"], "final_response": "今天北京晴,15-25度"} def call_schedule_tool(state: WorkflowState): # 模拟日程管理工具调用 return {"tool_calls": ["calendar_api"], "final_response": "已为您安排明天上午10点的会议"} def handle_general_query(state: WorkflowState): return {"final_response": f"您的问题「{state['user_input']}」已记录,我会尽快处理"} # 构建工作流 builder = StateGraph(WorkflowState) builder.add_node("classify", classify_task) builder.add_node("weather", call_weather_tool) builder.add_node("schedule", call_schedule_tool) builder.add_node("general", handle_general_query) builder.add_edge("classify", "weather") builder.add_edge("classify", "schedule") builder.add_edge("classify", "general") # 根据分类结果路由到不同节点 def route_after_classify(state: WorkflowState): task_type = state.get("task_type", "general_query") if task_type == "weather_query": return "weather" elif task_type == "schedule_manage": return "schedule" else: return "general" builder.add_conditional_edges( "classify", route_after_classify, { "weather": "weather", "schedule": "schedule", "general": "general" } ) builder.add_edge("weather", END) builder.add_edge("schedule", END) builder.add_edge("general", END) builder.set_entry_point("classify") workflow = builder.compile() # 测试不同输入 test_cases = [ "今天天气怎么样", "帮我安排明天的会议", "讲个笑话" ] for case in test_cases: result = workflow.invoke({"user_input": case}) print(f"输入: {case}") print(f"响应: {result['final_response']}") print("-" * 50)5.2 MCP 工具集成测试
接下来我们集成真实的 MCP 工具。以文件读取工具为例:
# mcp_integration.py import asyncio from mcp import ClientSession, StdioServerParameters from mcp.client.stdio import stdio_client async def test_mcp_tools(): # 配置 MCP Server(这里以标准输入输出为例) server_params = StdioServerParameters( command="python", args=["-m", "mcp.tools.filesystem"] # 假设有文件系统工具 ) async with stdio_client(server_params) as (read, write): async with ClientSession(read, write) as session: # 初始化会话 await session.initialize() # 列出可用工具 tools = await session.list_tools() print("可用工具:", tools) # 调用具体工具(示例) # result = await session.call_tool("read_file", {"path": "test.txt"}) # print("文件内容:", result) # 由于 MCP 工具需要具体实现,这里主要展示连接框架 # 实际使用时需要配置具体的 MCP Server5.3 完整智能小秘书演示
结合 LangGraph 工作流和 MCP 工具,构建完整的智能助手:
# smart_assistant.py from langgraph.graph import StateGraph, END from langgraph.prebuilt import create_react_agent from langchain_community.llms import Ollama # 或 OpenAI、Anthropic from typing import TypedDict, Annotated import operator class AssistantState(TypedDict): messages: Annotated[list, operator.add] current_step: str def setup_agent(): # 配置大模型(这里用 Ollama 本地模型示例) llm = Ollama(model="llama3:8b") # 定义工具列表(实际应连接 MCP Server) tools = [] # 创建 ReAct Agent agent = create_react_agent(llm, tools) return agent def main_workflow(state: AssistantState): agent = setup_agent() # 这里简化处理,实际应解析 agent 的思考过程 return {"messages": [{"role": "assistant", "content": "这是一个模拟响应"}]} # 构建完整工作流 builder = StateGraph(AssistantState) builder.add_node("assistant", main_workflow) builder.set_entry_point("assistant") builder.set_finish_point("assistant") assistant = builder.compile() # 测试对话 result = assistant.invoke({"messages": [{"role": "user", "content": "你好,请帮我查天气"}], "current_step": "start"}) print("助手响应:", result["messages"][-1]["content"])6. 接口 API 与批量任务
6.1 FastAPI 服务化
将智能小秘书封装成 HTTP API 服务:
# api_server.py from fastapi import FastAPI, HTTPException from pydantic import BaseModel import asyncio from smart_assistant import assistant # 引用之前的工作流 app = FastAPI(title="智能小秘书 API") class ChatRequest(BaseModel): message: str user_id: str = "default" class ChatResponse(BaseModel): response: str status: str @app.post("/chat", response_model=ChatResponse) async def chat_endpoint(request: ChatRequest): try: # 调用工作流 result = assistant.invoke({ "messages": [{"role": "user", "content": request.message}], "current_step": "start" }) return ChatResponse( response=result["messages"][-1]["content"], status="success" ) except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @app.post("/batch_chat") async def batch_chat_endpoint(requests: list[ChatRequest]): results = [] for req in requests: try: result = assistant.invoke({ "messages": [{"role": "user", "content": req.message}], "current_step": "start" }) results.append({ "user_id": req.user_id, "response": result["messages"][-1]["content"], "status": "success" }) except Exception as e: results.append({ "user_id": req.user_id, "response": "", "status": f"error: {str(e)}" }) return results if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)启动服务:
python api_server.py6.2 批量任务处理
对于需要处理大量任务的场景,可以结合队列实现:
# batch_processor.py import asyncio import json from concurrent.futures import ThreadPoolExecutor from smart_assistant import assistant class BatchProcessor: def __init__(self, max_workers=3): self.executor = ThreadPoolExecutor(max_workers=max_workers) def process_single_task(self, task_data): """处理单个任务""" try: result = assistant.invoke({ "messages": [{"role": "user", "content": task_data["message"]}], "current_step": "start" }) return { "task_id": task_data["id"], "status": "success", "response": result["messages"][-1]["content"] } except Exception as e: return { "task_id": task_data["id"], "status": "error", "error": str(e) } async def process_batch(self, tasks): """批量处理任务""" loop = asyncio.get_event_loop() futures = [ loop.run_in_executor(self.executor, self.process_single_task, task) for task in tasks ] results = await asyncio.gather(*futures) return results # 使用示例 async def demo_batch_processing(): processor = BatchProcessor() tasks = [ {"id": 1, "message": "今天天气如何"}, {"id": 2, "message": "安排明天会议"}, {"id": 3, "message": "查询新闻摘要"} ] results = await processor.process_batch(tasks) for result in results: print(f"任务 {result['task_id']}: {result['status']}") # 运行演示 asyncio.run(demo_batch_processing())7. 资源占用与性能观察
7.1 本地模型资源监控
如果你使用本地大模型(如通过 Ollama),需要关注以下资源指标:
显存占用观察:
# 监控 GPU 使用情况(需要 nvidia-smi) nvidia-smi -l 1 # 每秒刷新一次 # 或使用 gpustat(需安装:pip install gpustat) gpustat -i 1内存占用观察:
# 监控 Python 进程内存 ps aux | grep python | grep smart_assistant # 或使用 memory_profiler 进行详细分析 pip install memory_profiler python -m memory_profiler your_script.py7.2 API 调用性能优化
当使用云端 API 时,关注点转向网络延迟和费用控制:
# performance_monitor.py import time import asyncio from functools import wraps def api_metrics(func): """API 调用监控装饰器""" @wraps(func) async def wrapper(*args, **kwargs): start_time = time.time() try: result = await func(*args, **kwargs) duration = time.time() - start_time print(f"API 调用耗时: {duration:.2f}秒") return result except Exception as e: duration = time.time() - start_time print(f"API 调用失败,耗时: {duration:.2f}秒,错误: {e}") raise return wrapper # 在关键的模型调用函数上添加监控 @api_metrics async def call_llm_api(messages): # 模拟 API 调用 await asyncio.sleep(0.5) return "模拟响应"7.3 工作流状态持久化
对于长时间运行的工作流,需要实现状态保存和恢复:
# state_persistence.py import json import pickle from typing import Any class WorkflowStateManager: def __init__(self, storage_path="./workflow_states"): self.storage_path = storage_path def save_state(self, workflow_id: str, state: Any): """保存工作流状态""" filename = f"{self.storage_path}/{workflow_id}.pkl" with open(filename, 'wb') as f: pickle.dump(state, f) def load_state(self, workflow_id: str) -> Any: """加载工作流状态""" filename = f"{self.storage_path}/{workflow_id}.pkl" try: with open(filename, 'rb') as f: return pickle.load(f) except FileNotFoundError: return None def cleanup_state(self, workflow_id: str): """清理工作流状态""" filename = f"{self.storage_path}/{workflow_id}.pkl" try: os.remove(filename) except FileNotFoundError: pass # 使用示例 state_manager = WorkflowStateManager()8. 常见问题与排查方法
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 导入 LangGraph 报错 | 版本不兼容或安装不完整 | 检查 pip list 中的版本 | 使用 pip install langgraph==0.0.40 指定版本 |
| MCP 连接失败 | MCP Server 未启动或配置错误 | 检查 server_params 配置 | 确保 MCP Server 命令路径正确 |
| 工作流卡住不动 | 节点之间的边配置错误 | 打印状态流转日志 | 检查 add_edge 和 add_conditional_edges 配置 |
| 内存泄漏 | 工作流状态未及时清理 | 监控内存使用曲线 | 实现状态定期清理机制 |
| API 响应慢 | 网络延迟或模型负载高 | 添加超时和重试机制 | 使用异步调用和连接池 |
| 工具调用权限错误 | MCP Server 权限配置问题 | 检查工具调用参数 | 确保工具所需的访问权限 |
| 批量任务部分失败 | 单个任务异常影响整体 | 添加任务级别异常捕获 | 实现任务隔离和重试队列 |
8.1 依赖冲突解决
常见的依赖冲突及解决方法:
# 检查当前环境冲突 pip check # 如果发现冲突,创建干净环境重装 python -m venv clean_env source clean_env/bin/activate # 或 clean_env\Scripts\activate # 按顺序安装核心包 pip install langgraph==0.0.40 pip install langchain==0.1.0 pip install mcp==1.0.08.2 端口冲突处理
当启动多个服务时可能遇到端口占用:
# port_manager.py import socket from contextlib import closing def find_free_port(start_port=8000, end_port=9000): """查找可用端口""" for port in range(start_port, end_port + 1): with closing(socket.socket(socket.AF_INET, socket.SOCK_STREAM)) as sock: try: sock.bind(('localhost', port)) return port except socket.error: continue raise Exception("No free port found") # 在启动服务时使用 free_port = find_free_port() print(f"使用端口: {free_port}")9. 最佳实践与使用建议
9.1 开发阶段建议
循序渐进验证:
- 先从最简单的线性工作流开始,确保基础环境正常
- 逐步添加条件分支和循环控制
- 集成 MCP 工具时,先测试单个工具再组合使用
- 最后实现服务化和批量处理
代码组织规范:
# 推荐的项目结构 project/ ├── src/ │ ├── workflows/ # 工作流定义 │ ├── tools/ # MCP 工具集成 │ ├── models/ # 模型配置 │ └── api/ # API 服务 ├── tests/ # 测试用例 ├── config/ # 配置文件 └── requirements.txt # 依赖列表9.2 生产环境部署
安全配置:
# security_config.py from fastapi import Security from fastapi.security import APIKeyHeader api_key_header = APIKeyHeader(name="X-API-Key") def verify_api_key(api_key: str = Security(api_key_header)): """API 密钥验证""" valid_keys = ["your_secret_key_here"] # 从环境变量读取 if api_key not in valid_keys: raise HTTPException(status_code=403, detail="Invalid API Key") return api_key性能优化:
- 使用连接池管理数据库和外部服务连接
- 实现工作流结果的缓存机制
- 对频繁使用的工具调用结果进行缓存
- 设置合理的超时时间和重试策略
9.3 监控与日志
建立完整的监控体系:
# monitoring.py import logging import json from datetime import datetime # 配置结构化日志 logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s') logger = logging.getLogger("smart_assistant") def log_workflow_execution(workflow_id: str, user_input: str, response: str, duration: float): """记录工作流执行日志""" log_entry = { "timestamp": datetime.utcnow().isoformat(), "workflow_id": workflow_id, "user_input": user_input, "response": response, "duration_seconds": duration, "status": "success" if response else "error" } logger.info(json.dumps(log_entry))10. 总结与下一步
这个 LangGraph + MCP 的智能小秘书项目最值得尝试的点在于,它提供了一套完整的工作流编排方案,让复杂任务可以拆解成可管理的步骤。相比直接调用大模型,这种结构化的方式更容易调试、扩展和维护。
在实际部署时,建议先聚焦一个具体场景(比如天气查询+日程管理),把端到端的流程跑通,再逐步添加更多工具和能力。MCP 协议的优势在于工具的标准化的,这意味着你可以相对容易地替换或扩展工具集。
最容易遇到的坑主要是环境配置和依赖版本问题,特别是 LangGraph 和 MCP 都还在快速发展中,API 可能会有变动。解决方法是严格按照官方文档的版本要求,并且做好依赖隔离。
下一步可以探索的方向包括:集成更多类型的 MCP 工具(数据库、日历、邮件等)、实现工作流的热更新、添加用户反馈学习机制,或者将整个系统部署到云平台提供公开服务。
这个框架的扩展性很好,一旦掌握了核心模式,就可以快速构建出适应不同场景的智能助手。建议把本文的示例代码作为起点,根据实际需求进行调整和优化。