1. 为什么要把 MCP、A2A、Memory 放在一起做
如果你已经写过几个 LLM Agent Demo,大概率会遇到这样一类问题:工具调用写死在代码里,换一个模型或换一个工具就要重写一遍;多个 Agent 之间靠字符串拼接传递结果,稍微复杂一点就乱套;对话轮次一多,上下文窗口塞满,模型开始"失忆"。这三个问题分别对应 MCP、A2A、Memory 三个模块,单独看每个都不难,难的是让它们在同一个工程里协同跑通。
MCP(Model Context Protocol)解决的是"Agent 怎么标准化地调用外部工具和数据",你可以把它理解成 AI 世界的 USB-C 接口,工具方按协议暴露能力,Agent 方按协议发现和调用,双方不用互相认识。A2A(Agent-to-Agent)解决的是"多个 Agent 之间怎么分工协作",它更像一套组织流程,谁负责写代码、谁负责测试、谁负责审核,通过状态和消息在节点间流转。Memory 解决的是"Agent 怎么记住该记的东西",短期记忆保留最近几轮对话,长期记忆把历史压缩成摘要,避免上下文无限膨胀。
这篇教程面向已经能跑通基础 Agent 的开发者,目标是在本地把 MCP 服务端、A2A 消息路由、Memory 存储结构三块拼成一个可运行的链路。我会给出可复制的配置、代码骨架和验证步骤,源码放在文末仓库里。整条链路跑通后,你会得到一个能读文件、能翻译、能多 Agent 协作、还能记住上下文的完整 Agent 工程。
在开始之前,先说明一下模型接入的部分。本地调试时我用的是 TaoToken 提供的统一 API 入口,它兼容 OpenAI 的调用格式,省去了在不同模型厂商之间切换 SDK 的麻烦。官网地址是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,API 端点是 https://taotoken.net/api ,后面配置里会用到。
2. TaoToken 前置准备:拿到 Key 并确认模型可用
在写 MCP 服务端之前,先把模型调用这条线打通,否则后面调试工具调用时你分不清是协议问题还是模型问题。TaoToken 的接入方式和 OpenAI 完全一致,只需要替换 base_url 和 api_key。
第一步,登录控制台创建 API Key。打开 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite ,在 API Keys 页面新建一个密钥,复制保存。注意密钥只在创建时完整显示一次,丢了就重新建一个。
第二步,确认你要用的模型名称。在模型对话页面 https://taotoken.net/chat?utm_source=taotoken_aicg_blog_end&utm_content=model_chat&utm_campaign=rewrite 可以直接试聊,确认模型能正常响应,同时记下模型 ID,后面配置里要填。
第三步,把 Key 写进环境变量,不要硬编码在代码里。在项目根目录建一个.env文件:
# .env TAOTOKEN_API_KEY=sk-你的密钥 TAOTOKEN_BASE_URL=https://taotoken.net/api MODEL_NAME=你的模型ID然后用 python-dotenv 加载。这里有个小坑:.env文件一定要加进.gitignore,我见过有人把密钥提交到公开仓库,几分钟内就被扫走刷额度。
# config.py import os from dotenv import load_dotenv load_dotenv() API_KEY = os.getenv("TAOTOKEN_API_KEY") BASE_URL = os.getenv("TAOTOKEN_BASE_URL", "https://taotoken.net/api") MODEL_NAME = os.getenv("MODEL_NAME") if not API_KEY: raise RuntimeError("TAOTOKEN_API_KEY 未配置,请检查 .env 文件")第四步,写一个最小验证脚本,确认模型能调通:
# check_model.py from openai import OpenAI from config import API_KEY, BASE_URL, MODEL_NAME client = OpenAI(api_key=API_KEY, base_url=BASE_URL) resp = client.chat.completions.create( model=MODEL_NAME, messages=[{"role": "user", "content": "回复两个字:通了"}], ) print(resp.choices[0].message.content)运行python check_model.py,如果输出"通了",说明模型这条线没问题。如果报 401,检查 Key 是否复制完整;如果报 404,检查模型 ID 是否写对。这一步过了再往下走,能省掉后面大量排查时间。
3. MCP 服务端:用 JSON-RPC 暴露工具能力
MCP 的核心是 JSON-RPC 2.0 协议,Client 发请求,Server 执行并返回结果。我们先做一个最小可用的 MCP 服务端,暴露四个工具:两个数学运算、一个读 Markdown、一个写 Markdown。这样既能验证协议,又能为后面的翻译 Agent 提供文件读写能力。
3.1 项目结构与依赖
pip install fastapi uvicorn pydantic python-dotenv openai目录结构如下:
mcp-server/ ├── app/ │ ├── main.py │ ├── schemas.py │ ├── tools.py │ └── registry.py ├── data/ ├── requirements.txt └── run.sh3.2 定义 JSON-RPC 数据结构
MCP 的请求和响应都有固定字段,用 Pydantic 定义出来,既能做校验又能当文档:
# app/schemas.py from pydantic import BaseModel from typing import Any, Optional class JSONRPCRequest(BaseModel): jsonrpc: str id: Optional[int | str] = None method: str params: Optional[dict] = None class JSONRPCResponse(BaseModel): jsonrpc: str = "2.0" id: Optional[int | str] = None result: Any = None error: Optional[dict] = None这里jsonrpc固定为 "2.0",id用于匹配请求和响应,method是方法名,params是参数对象。响应里result和error二选一,成功返回 result,失败返回 error。
3.3 实现工具函数
工具函数就是普通的 Python 函数,参数和返回值都清晰定义。文件读写要限制在指定目录内,防止路径穿越:
# app/tools.py import os DOCUMENT_DIR = "./data" def add(a: int, b: int) -> int: return a + b def multiply(a: int, b: int) -> int: return a * b def read_markdown_tool(filename: str) -> dict: if not filename: raise ValueError("filename is required") if ".." in filename: raise ValueError("Invalid filename") file_path = os.path.join(DOCUMENT_DIR, filename) if not os.path.exists(file_path): raise ValueError("File not found") with open(file_path, "r", encoding="utf-8") as f: content = f.read() return {"filename": filename, "content": content} def write_markdown_tool(filename: str, content: str) -> dict: if not filename: raise ValueError("filename is required") if content is None: raise ValueError("content is required") if ".." in filename: raise ValueError("Invalid filename") file_path = os.path.join(DOCUMENT_DIR, filename) os.makedirs(DOCUMENT_DIR, exist_ok=True) with open(file_path, "w", encoding="utf-8") as f: f.write(content) return {"filename": filename, "status": "written", "length": len(content)}..检查是必须的,否则 Client 传../../etc/passwd就能读到系统文件。生产环境还要加白名单和权限校验。
3.4 工具注册表与调度层
把工具集中注册,方便tools/list动态返回:
# app/tools.py 追加 TOOLS = { "add": { "func": add, "description": "Add two integers", "parameters": {"a": "int", "b": "int"}, }, "multiply": { "func": multiply, "description": "Multiply two integers", "parameters": {"a": "int", "b": "int"}, }, "read_markdown_tool": { "func": read_markdown_tool, "description": "Read a markdown file", "parameters": {"filename": "string"}, }, "write_markdown_tool": { "func": write_markdown_tool, "description": "Write a markdown file", "parameters": {"filename": "string", "content": "string"}, }, }# app/registry.py from app.tools import TOOLS def list_tools(): return [ { "name": name, "description": meta["description"], "parameters": meta["parameters"], } for name, meta in TOOLS.items() ] def call_tool(name: str, params: dict): if name not in TOOLS: raise ValueError("Tool not found") func = TOOLS[name]["func"] return func(**params)3.5 主程序:暴露 /rpc 端点
# app/main.py from fastapi import FastAPI from app.schemas import JSONRPCRequest, JSONRPCResponse from app.registry import list_tools, call_tool app = FastAPI() @app.post("/rpc") async def rpc_handler(request: JSONRPCRequest): if request.jsonrpc != "2.0": return JSONRPCResponse( id=request.id, error={"code": -32600, "message": "Invalid JSON-RPC version"}, ) try: if request.method == "tools/list": result = list_tools() elif request.method == "tools/call": tool_name = request.params.get("name") arguments = request.params.get("arguments", {}) result = call_tool(tool_name, arguments) else: return JSONRPCResponse( id=request.id, error={"code": -32601, "message": "Method not found"}, ) return JSONRPCResponse(id=request.id, result=result) except Exception as e: return JSONRPCResponse( id=request.id, error={"code": -32000, "message": str(e)}, )启动服务:
uvicorn app.main:app --host 0.0.0.0 --port 8000到这里 MCP 服务端就跑起来了,它对外只暴露一个/rpc端点,所有能力通过 JSON-RPC 方法名区分。这种设计的好处是新增工具只需要改TOOLS字典,Client 端不用动。
4. A2A 消息路由:用状态图编排多 Agent 协作
MCP 解决的是单个 Agent 调用工具的问题,A2A 解决的是多个 Agent 之间怎么协作。我用 LangGraph 的状态图来实现,每个节点是一个 Agent,边是流转条件,共享状态在节点间传递。
4.1 定义共享状态
A2A 的核心是状态在 Agent 间流转,所以状态结构要设计清楚:
# a2a_workflow.py from typing import Annotated, TypedDict, List, Literal from langgraph.graph import StateGraph, END from langchain_core.messages import HumanMessage, AIMessage, BaseMessage class AgentState(TypedDict): messages: Annotated[List[BaseMessage], lambda x, y: x + y] current_stage: str task_status: Literal["pending", "success", "failed"]messages用Annotated加累加函数,保证每次节点返回的消息是追加而不是覆盖。task_status决定后续路由走向。
4.2 定义各司其职的 Agent 节点
每个 Agent 只做一件事,输入输出边界清晰:
def coder_agent(state: AgentState) -> dict: print("\n[Coder Agent] 正在分析需求并编写代码...") last_msg = state["messages"][-1].content if state["messages"] else "No input" generated_code = f"def solve_task():\n # Solving: {last_msg}\n return 'Success'" return { "messages": [AIMessage(content=f"Code Generated:\n{generated_code}")], "current_stage": "testing", "task_status": "pending", } def tester_agent(state: AgentState) -> dict: print("\n[Tester Agent] 正在运行单元测试和边界检查...") test_result = "All tests passed. Coverage: 98%." status = "success" return { "messages": [AIMessage(content=f"Test Report: {test_result}")], "current_stage": "review" if status == "success" else "fixing", "task_status": status, } def fixer_agent(state: AgentState) -> dict: print("\n[Fixer Agent] 检测到失败,正在紧急修复...") return { "messages": [AIMessage(content="Bug fixed. Re-submitting for testing.")], "current_stage": "testing", "task_status": "pending", } def reviewer_agent(state: AgentState) -> dict: print("\n[Reviewer Agent] 进行最终代码审查...") return { "messages": [AIMessage(content="Code Review Passed. Ready for deployment.")], "current_stage": "done", "task_status": "success", }4.3 编排状态图与条件路由
workflow = StateGraph(AgentState) workflow.add_node("coder", coder_agent) workflow.add_node("tester", tester_agent) workflow.add_node("fixer", fixer_agent) workflow.add_node("reviewer", reviewer_agent) workflow.set_entry_point("coder") workflow.add_edge("coder", "tester") def route_after_test(state: AgentState) -> Literal["fixer", "reviewer"]: if state["task_status"] == "failed": return "fixer" return "reviewer" workflow.add_conditional_edges( "tester", route_after_test, {"fixer": "fixer", "reviewer": "reviewer"}, ) workflow.add_edge("fixer", "tester") workflow.add_edge("reviewer", END) app = workflow.compile()这里的关键是route_after_test,它根据task_status决定测试失败后走修复还是走审核。修复完再回到测试,形成闭环。这种结构比线性调用灵活得多,新增一个 Agent 只需要加节点和边。
4.4 运行验证
if __name__ == "__main__": initial_input = { "messages": [HumanMessage(content="Create a function to calculate Fibonacci sequence")], "current_stage": "coding", "task_status": "pending", } final_state = app.invoke(initial_input) for msg in final_state["messages"]: role = "User" if isinstance(msg, HumanMessage) else "Agent" print(f"[{role}]: {msg.content[:100]}")运行后会看到 Coder → Tester → Reviewer 依次执行,如果测试失败会自动插入 Fixer 再回到 Tester。这套骨架可以直接套用到代码生成、数据分析、客服工单等场景。
5. Memory 存储结构:短期记忆加长期摘要
Memory 要解决的是上下文窗口有限的问题。我的做法是短期记忆保留最近几轮完整对话,超出阈值后把最老的一轮压缩成摘要,追加到长期记忆里。
5.1 记忆结构设计
# memory_agent.py from openai import OpenAI from config import API_KEY, BASE_URL, MODEL_NAME client = OpenAI(api_key=API_KEY, base_url=BASE_URL) class CompressedAgentMemory: def __init__(self, model=MODEL_NAME): self.model = model self.short_term_memory = [] self.long_term_summary = "" self.max_short_term_rounds = 3 def _summarize_memory(self, old_summary, new_memory): prompt = f"""你是一个记忆管理总结专家,请将以下对话内容浓缩并更新到现有的背景摘要中。 要求:保留关键事实(姓名、职业、兴趣、决定、事件等)。 现有背景摘要:{old_summary} 新的对话内容:{new_memory} 请生成一个新的、内容精炼的背景摘要。""" response = client.chat.completions.create( model=self.model, messages=[{"role": "user", "content": prompt}], ) return response.choices[0].message.content短期记忆是一个列表,存 user 和 assistant 的消息对。长期记忆是一个字符串,存压缩后的摘要。阈值max_short_term_rounds控制什么时候触发压缩。
5.2 对话与压缩逻辑
def chat(self, user_input): system_prompt = f"""你是一个智能助手。 以下是之前的对话背景摘要(长期记忆): {self.long_term_summary if self.long_term_summary else "无长期记忆"}""" messages = [{"role": "system", "content": system_prompt}] messages.extend(self.short_term_memory) messages.append({"role": "user", "content": user_input}) response = client.chat.completions.create( model=self.model, messages=messages, ) assistant_reply = response.choices[0].message.content self.short_term_memory.append({"role": "user", "content": user_input}) self.short_term_memory.append({"role": "assistant", "content": assistant_reply}) if len(self.short_term_memory) > self.max_short_term_rounds * 2: need_compress = self.short_term_memory[:2] self.short_term_memory = self.short_term_memory[2:] compress_str = ( f"User: {need_compress[0]['content']}\n" f"Assistant: {need_compress[1]['content']}" ) self.long_term_summary = self._summarize_memory( self.long_term_summary, compress_str ) return assistant_reply每次对话后检查短期记忆长度,超过阈值就把最老的一轮拿出来压缩。压缩后的摘要替换掉原来的长期记忆,短期记忆里删掉那一轮。这样上下文长度始终可控,同时关键信息不会丢。
5.3 验证记忆效果
if __name__ == "__main__": agent = CompressedAgentMemory() print("Round1:", agent.chat("你好我是孙海亮,我是一名全栈AI工程师")) print("Round2:", agent.chat("我最近在研究AI LLM技术")) print("Round3:", agent.chat("我想开发一个帮我写代码的工具!")) print("Round4:", agent.chat("我是谁?")) print("Round5:", agent.chat("我最近在研究什么技术?"))跑下来你会发现,即使前几轮已经被压缩进长期摘要,Round4 和 Round5 依然能正确回答"你是谁"和"你在研究什么"。这就是短期加长期双层记忆的价值。
6. 三模块协同:把 MCP、A2A、Memory 串成一条链路
前面三块单独都能跑,现在把它们拼起来。我做一个 Markdown 翻译 Agent:MCP 提供文件读写工具,Memory 管理翻译进度,A2A 负责在读取、拆分、翻译、写入几个阶段之间流转。
6.1 MCP 客户端封装
# mcp_http_client.py import requests class MCPHttpClient: def __init__(self, url: str): self.url = url self.session = requests.Session() def call(self, method: str, params: dict = None, request_id: int = 1): payload = { "jsonrpc": "2.0", "id": request_id, "method": method, "params": params or {}, } response = self.session.post(self.url, json=payload, timeout=5) response.raise_for_status() data = response.json() err = data.get("error") if err is not None: raise RuntimeError(f"RPC Error: {err}") return data.get("result") def list_tools(self): return self.call("tools/list") def call_tool(self, name: str, arguments: dict): return self.call("tools/call", {"name": name, "arguments": arguments})6.2 上下文记忆结构
翻译场景的记忆和对话记忆不同,它要记录原文、分片、游标和输出缓冲:
# context_memory.py from typing import List class ContextMemory: def __init__(self, model_name: str): self.row_content = "" self.source_chunks: List[str] = [] self.read_cursor: int = 0 self.output_buffer: List[str] = [] self.model_name = model_name def get_full_output(self) -> str: return "\n\n".join(self.output_buffer)6.3 本地工具:拆分与翻译
# local_tools.py from langchain_text_splitters import RecursiveCharacterTextSplitter def local_tool_split(memory): splitter = RecursiveCharacterTextSplitter(chunk_size=800, chunk_overlap=100) chunks = splitter.split_text(memory.row_content) memory.source_chunks = chunks memory.read_cursor = 0 memory.output_buffer = [] return f"拆分完成,内存已有{len(chunks)}个chunk" def local_tool_get_next_chunk_and_translate(memory, model_client): if memory.read_cursor >= len(memory.source_chunks): return "EOF" chunk = memory.source_chunks[memory.read_cursor] idx = memory.read_cursor memory.read_cursor += 1 response = model_client.chat.completions.create( model=memory.model_name, messages=[ {"role": "system", "content": "直译Markdown为中文,保留格式。"}, {"role": "user", "content": chunk}, ], temperature=0.1, ) return response.choices[0].message.content def local_tool_save(content, memory): memory.output_buffer.append(content) return f"已保存,当前缓冲区{len(memory.output_buffer)}个chunk"6.4 主循环:工具调用与历史裁剪
主循环的逻辑是:模型决定调用哪个工具,我们执行工具,把结果反馈给模型,直到模型不再调用工具。关键点是每次工具调用后裁剪历史,只保留 system、首个 user 和最新状态,避免 token 爆炸。
# run_agent.py import json from openai import OpenAI from config import API_KEY, BASE_URL, MODEL_NAME from mcp_http_client import MCPHttpClient from context_memory import ContextMemory from local_tools import ( local_tool_split, local_tool_get_next_chunk_and_translate, local_tool_save, ) client = OpenAI(api_key=API_KEY, base_url=BASE_URL) memory = ContextMemory(MODEL_NAME) LOCAL_TOOL_SCHEMA = [ { "type": "function", "function": { "name": "local_tool_split", "description": "将内存中的长文本拆分成片段", "parameters": {"type": "object", "properties": {}}, }, }, { "type": "function", "function": { "name": "local_tool_get_next_chunk_and_translate", "description": "取出下一段翻译,返回EOF表示完成", "parameters": {"type": "object", "properties": {}}, }, }, ] def run_agent(): mcp_client = MCPHttpClient("http://127.0.0.1:8000/rpc") mcp_tools = mcp_client.list_tools() mcp_tools_map = {tool["name"]: tool for tool in mcp_tools} openai_tools = [ { "type": "function", "function": { "name": tool["name"], "description": tool["description"], "parameters": tool["parameters"], }, } for tool in mcp_tools ] all_tools = openai_tools + LOCAL_TOOL_SCHEMA messages = [ { "role": "system", "content": ( "你是一个严谨的 Markdown 翻译 Agent。" "先调用 read_markdown_tool 读取文件," "再调用 local_tool_split 拆分," "然后循环调用 local_tool_get_next_chunk_and_translate 直到返回 EOF," "最后调用 write_markdown_tool 写入结果。" ), }, {"role": "user", "content": "请读取 input.md,翻译为中文并保存到 output_zh.md。"}, ] should_prune = False while True: response = client.chat.completions.create( model=MODEL_NAME, messages=messages, tools=all_tools, tool_choice="auto", ) assistant_message = response.choices[0].message messages.append(assistant_message) if not assistant_message.tool_calls: break if len(assistant_message.tool_calls) > 1: assistant_message.tool_calls = [assistant_message.tool_calls[0]] for tool_call in assistant_message.tool_calls: name = tool_call.function.name raw_args = tool_call.function.arguments args = json.loads(raw_args) if isinstance(raw_args, str) else (raw_args or {}) tool_result = "" if name in mcp_tools_map: if name == "write_markdown_tool" and args.get("content") == "__USE_MEMORY_BUFFER__": args["content"] = memory.get_full_output() res = mcp_client.call_tool(name, args) if name == "read_markdown_tool": memory.row_content = res.get("content") tool_result = f"文件读取成功,长度{len(memory.row_content)},请调用 local_tool_split" else: tool_result = "执行成功(MCP TOOL)" elif name == "local_tool_split": local_tool_split(memory) tool_result = "[已省略大文本]" should_prune = True elif name == "local_tool_get_next_chunk_and_translate": result = local_tool_get_next_chunk_and_translate(memory, client) if result != "EOF": local_tool_save(result, memory) tool_result = f"第{memory.read_cursor}段完成" else: tool_result = "EOF" should_prune = True messages.append({ "role": "tool", "tool_call_id": tool_call.id, "content": str(tool_result), }) if should_prune: messages = [ messages[0], messages[1], {"role": "user", "content": "系统通知:继续下一步"}, ] should_prune = False if __name__ == "__main__": run_agent()这段代码把三个模块串起来了:MCP 负责文件读写,Memory 负责进度和缓冲,主循环负责工具调度和历史裁剪。跑通后你会看到模型自动完成读取、拆分、逐段翻译、写入的完整流程。
7. 本篇常见错排查
MCP 服务端启动报 422 错误。多半是请求体不符合JSONRPCRequest的字段定义。检查jsonrpc是否为 "2.0",method是否拼写正确。用 curl 直接测:
curl -X POST http://127.0.0.1:8000/rpc \ -H "Content-Type: application/json" \ -d '{"jsonrpc":"2.0","id":1,"method":"tools/list","params":{}}'工具调用返回 "Tool not found"。检查TOOLS字典里的 key 和模型传的name是否一致。模型有时会自己编工具名,所以 system prompt 里要把可用工具列清楚。
A2A 状态图报 KeyError。检查AgentState里定义的字段是否和节点返回的 dict key 一致。LangGraph 对状态字段是严格匹配的,多一个少一个都会报错。
Memory 压缩后模型失忆。检查_summarize_memory的 prompt 是否明确要求保留关键事实。摘要太简略会丢信息,可以在 prompt 里加"保留姓名、职业、决定、事件"这类约束。
翻译循环停不下来。检查local_tool_get_next_chunk_and_translate的 EOF 判断,游标是否正常递增。如果游标没动,会一直返回同一段。
历史裁剪后模型重复调用工具。裁剪时保留了 system 和首个 user,但模型可能忘了当前进度。解决办法是在裁剪后的 user 消息里带上进度摘要,比如"已完成第 3/10 段"。
API 调用报 401 或 404。401 检查 Key 是否完整,404 检查模型 ID 和 base_url。TaoToken 的 base_url 是 https://taotoken.net/api ,不要漏掉/api。
8. 下一步:把链路接到你的项目里
到这里,MCP 服务端、A2A 状态图、Memory 双层记忆三块都跑通了,而且拼成了一条完整的翻译 Agent 链路。你可以把这套结构直接搬到自己的项目里:把TOOLS字典换成你的业务工具,把 A2A 的节点换成你的业务 Agent,把 Memory 的压缩策略换成适合你场景的阈值。
如果你在接入模型时想统一管理多个厂商的 Key,可以在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api_keys&utm_campaign=rewrite 创建密钥,接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 有完整的参数说明。长期做编码类 Agent 的话,Coding Plan 页面 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding_plan&utm_campaign=rewrite 有更细的配额方案,适合需要频繁调用模型的场景。
源码仓库在 https://github.com/sunhailiang-tourist/agent-and-mcp ,包含本篇所有代码和示例数据。建议先按第 2 节验证模型,再按第 3 节起 MCP 服务,最后跑第 6 节的主循环。每一步都单独验证过再往下走,比一次性全跑起来再排查要快得多。