基于A2A协议的智能体协作系统——SmartVoyage智途v7.8架构设计
一、项目背景
SmartVoyage智途是一个面向旅游场景的AI智能体协作系统,支持机票查询、火车票预订、演唱会门票、租车、保险等丰富功能。核心创新在于采用了A2A(Agent-to-Agent)协议实现多Agent之间的标准化通信,让每个Agent专注自己的垂直领域,通过ReAct范式和依赖分组并行机制自动完成复杂的多意图任务规划与执行。
技术栈:
Python 3.10 + FastAPI
A2A协议 (Agent-to-Agent)
MCP协议 (Model Context Protocol)
Docker容器化部署
Redis精确缓存 + Milvus语义缓存
MySQL业务数据库 + Conversation Memory
OpenAI/兼容API (ChatOpenAI)
二、整体架构图
以下为系统六层七服务的分层架构全景:

系统采用六层七服务分层架构,每一层职责清晰、松耦合。下面逐层展开。
第1层:API网关层
# main_http_api.py — FastAPI网关
app = FastAPI(title="SmartVoyage API")
@app.post("/api/chat") # 同步接口
@app.post("/api/chat/stream") # SSE流式接口
@app.post("/api/chat_async") # 异步接口
@app.get("/api/memory/*") # 记忆操作
@app.get("/api/agents") # Agent列表
@app.get("/health") # 健康检查
- uvicorn监听 8888端口
- 提供同步/流式/异步三种API接口,适配不同客户端
- 双层缓存拦截(Redis→Milvus),命中直接返回,未命中才进入推理链路
- 这是唯一暴露给外部的入口,内部服务端口全部隔离
第2层:业务入口层
# main.py — ChatService门面类
class ChatService:
def __init__(self, user_id: str): ...
async def get_agent_cards() -> list[dict] # 获取所有Agent能力卡片
async def update_user_profile(profile: dict) # 更新用户画像
async def get_memory_state() # 读取四种记忆
async def clear_memory() # 清除记忆
async def chat(user_input: str) -> str # ★ 核心对话入口
采用门面模式统一对外接口,隐藏内部复杂性。每个user_id对应一个独立的ChatService实例,实现多用户隔离。
第3层:推理核心层(★系统大脑)
# core/chat_service_core.py — VoyageChatCore
class VoyageChatCore:
def __init__(self, user_id: str):
self.agent_network = AgentNetwork(...) # Agent网络
self.memory = ConversationMemory(...) # 加载记忆
self.cache_manager = CacheManager(...) # 加载缓存
self.llm = ChatOpenAI(...) # 全局单例LLM
# 推理五步曲
def _intent_agent(query) -> IntentTriple # 意图识别
def _plan_or_reflect(steps, observations) # 任务规划/反思
def _react_loop(steps) # ReAct循环
async def chat(user_input) -> str # 主入口
这是整个系统的核心,实现了完整的推理链路:
- 意图识别 → 判断用户有几个意图
- 任务规划 → 多意图时生成步骤计划
- ReAct循环 → 分组并行执行 + 反思调整
- 结果汇总 → 合并多个Agent的返回
第4层:A2A Server层 (:6001-6003)
三个独立的FastAPI应用,分别实现旅游领域的专业Agent:

每个A2A Server实现handle_task(task)方法,遵循TaskState状态机:
Started → InProgress → Completed / Acknowledged
第5层:MCP Server层 (:5001-5003)
MCP (Model Context Protocol) 是标准化的工具调用协议,每个Server包装一组业务函数:
# mcp_server/mcp_ticket.py
server = FastMCP("ticket_mcp", host="0.0.0.0", port=5002)
@server.tool()
def query_train_tickets(from_city, to_city, date): ...
@server.tool()
def order_train_ticket(order_data): ...
- 共6个@mcp.tool装饰器函数(火车票/航班/演唱会各2个:query + order)
- 使用Streamable HTTP传输协议
- uvicorn启动,注册到Docker网络
第6层:mcp_tools层(具体业务逻辑)
实际执行SQL查询的模块:
# mcp_tools/ticket_tools.py
def query_train_tickets(from_city, to_city, date):
sql = f"SELECT * FROM trains WHERE departure='{from_city}'..."
return mysql.execute(sql)
def order_train_ticket(order_data):
sql = f"INSERT INTO orders ... VALUES (...)"
return mysql.execute(sql)
共涉及6个查询函数 + 6个下单函数 + 天气采集函数,直连MySQL。
三、横向支撑组件
1. 四种记忆系统
# core/memory_manager.py
class ConversationMemory:
MEMORY_TYPES = {
"short_term": {"max_messages": 10}, # 最近10条对话
"user_profile": {}, # 用户偏好画像
"task_context": {}, # 当前任务上下文
"long_term": {"max_entries": 20}, # 长期关键信息
}
- 每种记忆独立管理,互不干扰
- add_message(role, content) 写入
- get_memory(type, limit) 读取
- 所有非临时记忆持久化到MySQL的connection_memory表
- 通过user_id字段实现多用户隔离
2. 双层缓存机制
# core/cache_manager.py
class CacheManager:
def check_cache(user_input) -> str | None:
# L1: Redis精确匹配
md5_key = hashlib.md5(user_input.encode()).hexdigest()
cached = redis.get(md5_key)
if cached:
return {"text": cached, "cached": True}
# L2: Milvus语义匹配
embedding = milvus.generate_embedding(user_input)
results = milvus.search(embedding, top_k=1, threshold=0.92)
if results and results[0].distance >= 0.92:
return {"text": results[0].payload, "cached": True, "milvus": True}
return None
- L1 Redis: 基于MD5(key),命中率最高的快速通道
- L2 Milvus: 基于COSINE余弦相似度≥0.92,覆盖同义表达
- Cache Promotion: 当Milvus命中之后,自动创建Redis键值对,后续相同含义的请求走Redis
四、核心数据流:以"帮我查明天北京到上海的机票和上海天气"为例
一轮完整对话的决策树流程如下(含缓存拦截、意图识别、ReAct循环、记忆更新):
flowchart TD
Start([用户输入 query]) --> CacheCheck{双层缓存命中?}
CacheCheck -- Redis L1 命中 --> CacheHitL1[返回 cached_result<br>更新 stats]
CacheCheck -- Milvus L2 命中 --> CacheHitL2[返回 result + milvus=true<br>创建 Redis key - Promotion]
CacheCheck -- 未命中 --> Gateway[API网关接收<br>HTTP POST /api/chat]
CacheHitL1 --> ReturnEnd
CacheHitL2 --> Gateway
Gateway --> ChatSvc[ChatService.chat()]
ChatSvc --> LoadMem[(加载四种记忆)]
LoadMem --> IntentAgent[意图识别 Agent<br>注入7个变量到 Prompt]
IntentAgent --> IntentOutput{LLM输出三元组}
IntentOutput -->|follow_up不为空| FollowUpResp["有追问!<br>直接返回 follow_up_message"]
FollowUpResp --> MemUpdate
IntentOutput -->|intents为空| NoIntent["无有效意图<br>回复请重新描述"]
NoIntent --> ReturnEnd
intents_num=intents数量 > 1 ?)
NoIntent --> ReturnEnd
FollowUpResp --> intents_num
IntentOutput -->|单一意图| SingleInt{== 单一意图 ==}
SingleInt --> YesSingle["是 → _run_single_intent()"]
YesSingle --> FindAgent[查找 intent_to_agent_map<br>定位对应A2A Server]
FindAgent --> A2CCall[A2AClient.send_task_async<br>RPC调用A2A Server]
A2CCall --> MCPExec[MCP Server<br>@mcp.tool执行SQL]
MCPExec --> Summarize[LLM总结原始JSON]
Summarize --> MemUpdate
SingleInt --否 (多意图)--> MultiInt["是 → _plan_or_reflect()<br>生成步骤计划"]
MultiInt --> ReactLoop{"ReAct循环<br>max_rounds=3"}
ReactLoop --> GroupPlan[按depends_on分组<br>OrderedDict排序]
GroupPlan --> AsyncGather[asyncio.gather()组内并行]
AsyncGather --> CollectAll[收集 all_results]
CollectAll --> HasValid{有有效结果?}
HasValid -- 是 --> CollectAll
HasValid -- 否 --> Reflect[观察结果拼接为obs_text<br>再次调用 _plan_or_reflect<br>判断是否需要新步骤]
Reflect --> NewSteps{"有新步骤?<br>及 max_rounds未到?"}
NewSteps -- 是 --> ReactLoop
NewSteps -- 否 --> BreakLoop[break ReAct循环]
BreakLoop --> ResultCount[检查结果数量]
CollectAll --> ResultCount
ResultCount --> CountOne{== 仅1个结果 ==}
CountOne -- 是 --> DirectReturn["直接返回该结果"]
CountOne -- 否 (≥2)--> MergeResult["调用 react_summary_prompt<br>LLM合并多结果为自然语言"]
MergeResult --> MemUpdate
DirectReturn --> MemUpdate
MemUpdate[memory.add_message<br>cache_manager.write_cache<br>写入Redis + Milvus] --> ReturnEnd([返回给用户])
style Start fill:#74b9ff,color:#fff,stroke:#0984e3
style CacheCheck fill:#55efc4,color:#000
style ReactLoop fill:#fdcb6e,color:#000,stroke:#f39c12
style Intents_output fill:#ffeaa7
style MemUpdate fill:#dfe6e9,stroke:#636e72
style ReturnEnd fill:#74b9ff,color:#fff
style SingleInt fill:#fab1a0
style MultiInt fill:#fab1a0
style CollectAll fill:#81ecec
完整流程拆解为8步:
Step 1: 用户发送消息 → API网关接收
HTTP POST /api/chat:
{
"user_id": "user001",
"query": "帮我查明天北京到上海的机票和上海天气"
}
Step 2: 双层缓存检查
- Redis MD5查询 → 未命中
- Milvus语义搜索 → 未命中
- 继续进入推理链路
Step 3: 意图识别(Intent Agent)
注入7个变量到prompt:
- available_tools — 可用工具列表
- intent_info — 意图类型说明
- conversation_history — 历史对话
- query — 当前问题
- current_date — 当前日期
- user_profile — 用户画像
- task_context — 任务上下文
LLM输出三元组:
{
"intents": ["flight_query", "weather_query"],
"user_queries": {
"flight_query": "明天北京到上海的机票",
"weather_query": "上海明天的天气"
},
"follow_up_message": null
}
三分支判断:
- intents数量 = 2 > 1 → 进入多意图规划流程
Step 4: 任务规划(Plan)
首次调用_plan_or_reflect(observations=''),LLM输出步骤计划:
{
"steps": [
{"step": 1, "action": "query_flight", "intent": "flight_query", "depends_on": [0]},
{"step": 2, "action": "query_weather", "intent": "weather_query", "depends_on": [0]}
]
}
两个步骤互相独立(都depends_on [0]表示无依赖),可以并行执行。
Step 5: ReAct循环 — 分组并行执行
按depends_on排序,使用OrderedDict:
Group 1 (indep): Step1(航班), Step2(天气) → asyncio.gather() 并行
每个step走到_run_single_intent:
- 查找intent_to_agent_map["flight_query"] → 定位到a2a-ticket:6002
- A2AClient.send_task_async(task, url="http://a2a-ticket:6002")
- HTTP POST到A2A Server → TaskResponse
- A2A内部调用MCP Server:5002 → @mcp.tool(query_flight) → SQL查询
- 原始JSON → config.yaml总结提示词 → LLM总结为自然语言
两个分支同时完成,结果收集到all_results = [result1, result2]。
Step 6: 反思(Reflect)
检查结果数量 = 2,有有效结果。
因为两个步骤独立(depends_on全是0),跳过反思直接break。
如果不是全独立,则拼接obs_text = "步骤1结果...\n步骤2结果...",再次调用_plan_or_reflect(observations=obs_text),看是否有新步骤产生。最多3轮。
Step 7: 结果汇总
结果数量 = 2 → 多结果合并
调用react_summary_prompt注入所有结果和原始query:
"请根据以下信息回答用户问题..."
LLM输出一段连贯的自然语言回复。
Step 8: 记忆更新 + 缓存写入
memory.add_message("user", "帮我查明天北京到上海的机票...")
memory.add_message("assistant", "为您找到以下航班...")
cache_manager.write_cache(user_input, final_result)
最后返回给用户。
五、Docker部署架构
以下为11个容器的编排拓扑与三层代理设计:
graph LR
subgraph INTERNET ["互联网"]
USER[用户浏览器 / 客户端]
end
subgraph LAYER1_Nginx ["第一层:反向代理 - 唯一对外端口"]
NGINX_EXT["nginx-proxy | port: 8888 HTTPS | → nginx:80"]
end
subgraph LAYER2_Gateway ["第二层:网关/服务层"]
APIGW["api-gateway :8888 | FastAPI + uvicorn | 推理核心"]
subgraph MCP_PROXY_NET ["MCP路由代理集群"]
MCPPROXY["nginx:mcp-proxy | 5001-5003"]
end
subgraph A2A_PROXY_NET ["A2A路由代理集群"]
A2APROXY["nginx:a2a-proxy | 6001-6003"]
end
end
subgraph MCP_SERVICES ["第三层:MCP Server (:5001-5003)"]
MCP_WEATHER["mcp-weather :5001 | @mcp.tool(x6)]
MCP_TICKET["mcp-ticket :5002 | @mcp.tool(x6)]
MCP_TRIP["mcp-trip :5003 | @mcp.tool(x6)]
end
subgraph A2A_SERVICES ["第四层:A2A Agent (:6001-6003)"]
A2A_WEATHER["a2a-weather :6001 | AgentCard + handle_task"]
A2A_TICKET["a2a-ticket :6002 | AgentCard + handle_task"]
A2A_TRIP["a2a-trip :6003 | AgentCard + handle_task"]
end
subgraph COLLECTOR["后台进程"]
WCOLLECT[weather-collector<br>定时采集天气数据]
end
subgraph STORAGE["持久化存储 (宿主机/容器外)"]
MYSQL_DB[(MySQL DB - 业务数据 + 记忆表)]
MILVUS_DB[(Milvus DB - 向量数据库)]
REDIS_DB[(Redis - 键值缓存)]
end
USER -->|HTTPS:8888| NGINX_EXT
NGINX_EXT --> api:[FastAPI endpoints]
NGINX_EXT --> mcp:[MCP routes :5001-5003]
NGINX_EXT --> a2a:[A2A routes :6001-6003]
api --- APIGW
mcp --- MCPPROXY
a2a --- A2APROXY
MCPPROXY --> MCP_WEATHER
MCPPROXY --> MCP_TICKET
MCPPROXY --> MCP_TRIP
A2APROXY --> A2A_WEATHER
A2APROXY --> A2A_TICKET
A2APROXY --> A2A_TRIP
A2A_WEATHER --> MCP_WEATHER
A2A_TICKET --> MCP_TICKET
A2A_TRIP --> MCP_TRIP
MCP_WEATHER -.-> MYSQL_DB
MCP_TICKET -.-> MYSQL_DB
MCP_TRIP -.-> MYSQL_DB
WCOLLECT -->|写入天气数据| MYSQL_DB
APIGW -.-> REDIS_DB
APIGW -.-> MILVUS_DB
style USER fill:#74b9ff,stroke:#0984e3,color:#fff
style NGINX_EXT fill:#fdcb6e,stroke:#f39c12,stroke-width:3px,color:#000
style APIGW fill:#74b9ff,stroke:#0984e3
style MCPPROXY fill:#b2bec3,stroke:#636e72
style A2APROXY fill:#b2bec3,stroke:#636e72
style MCP_WEATHER fill:#a29bfe,stroke:#6c5ce7
style MCP_TICKET fill:#a29bfe,stroke:#6c5ce7
style MCP_TRIP fill:#a29bfe,stroke:#6c5ce7
style A2A_WEATHER fill:#dfe6e9,stroke:#636e72
style A2A_TICKET fill:#dfe6e9,stroke:#636e72
style A2A_TRIP fill:#dfe6e9,stroke:#636e72
style MYSQL_DB fill:#74b9ff,stroke:#0984e3
style MILVUS_DB fill:#a29bfe,stroke:#6c5ce7
style REDIS_DB fill:#ff6b6b,stroke:#ee5a24
容器清单
| 容器名 | 镜像 | 端口 | 说明 |
|---|---|---|---|
| nginx-proxy | nginx:latest | 8888→80 | 唯一对外入口 |
| api-gateway | python-api | 8888 | FastAPI推理核心 |
| mcp-proxy | nginx:latest | 5001-5003 | MCP路由代理 |
| a2a-proxy | nginx:latest | 6001-6003 | A2A路由代理 |
| mcp-weather | python-mcp | 5001 | 天气工具Server |
| mcp-ticket | python-mcp | 5002 | 票务工具Server |
| mcp-trip | python-mcp | 5003 | 行程工具Server |
| a2a-weather | python-a2a | 6001 | 天气Agent |
| a2a-ticket | python-a2a | 6002 | 票务Agent |
| a2a-trip | python-a2a | 6003 | 行程Agent |
| weather-collector | python-collect | - | 定时采集脚本 |
三层代理设计
互联网 → nginx:8888 → api-gateway:8888
↓
mcp-proxy(5001-5003) / a2a-proxy(6001-6003)
↓
各微服务容器 (Docker内部网络)
外部只能访问8888端口,内部通过Docker Network互相调用,安全性极高。
六、设计亮点总结
1. 配置驱动动态注册
加一个新Agent只需改config.yaml,无需修改代码:
agents:
- name: hotel_agent
a2a_url: http://a2a-hotel:6004
mcp_url: http://mcp-hotel:5004
重启后自动生效,符合开闭原则。
2. 依赖分组并行执行
asyncio.gather(*tasks) # 组内并发
# 下一组
asyncio.gather(*tasks) # 组间串行
最大化利用并行能力,有依赖的步骤不会提前执行。
3. ReAct反思循环
规划(仅步骤) → 执行 → 观察结果 → 再规划(带observations)
↓
新步骤为空? → break
↓ 不是
max_rounds <= 3?
最多3轮"规划→执行→反思→再规划",避免无限循环。
4. 双层缓存Promotion
用户问:"北京明天天气怎么样"
→ Milvus语义命中(之前有人问过类似表述)
→ 返回结果 + cached=True
→ 自动创建Redis key: MD5("北京明天天气怎么样")
→ 下次同样字面请求直接走Redis
语义缓存充当桥梁,逐步积累精确缓存命中率。
5. 四种记忆融合
短中期对话保持连贯性,用户偏好塑造个性化体验,任务上下文保证多轮复杂对话不出错,长期记忆沉淀有价值信息。
七、总结
SmartVoyage智途将A2A协议、MCP协议、ReAct范式有机结合,实现了可扩展的智能体协作架构。通过配置驱动和分层设计,新增功能模块成本极低。核心难点在于意图识别的Prompt设计、ReAct循环的终止条件判定、以及多Agent并行执行的结果合并。
下一篇文章将详细拆解ReAct推理链路的代码实现和Prompt工程技巧,欢迎关注。
更多推荐


所有评论(0)