一、项目背景


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              # 主入口

这是整个系统的核心,实现了完整的推理链路:

  1. 意图识别 → 判断用户有几个意图
  2. 任务规划 → 多意图时生成步骤计划
  3. ReAct循环 → 分组并行执行 + 反思调整
  4. 结果汇总 → 合并多个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工程技巧,欢迎关注。

Logo

欢迎加入 MCP 技术社区!与志同道合者携手前行,一同解锁 MCP 技术的无限可能!

更多推荐