跨境AI智能体分层架构设计:网关层、调度层与执行层的解耦实践
一、为什么需要分层解耦
在跨境电商场景中,AI智能体需要同时处理客服对话、广告投放、供应链查询等多种业务,背后涉及OpenAI、Claude、DeepSeek等多个大模型厂商的API调用,以及Amazon、Walmart、ERP等异构系统的对接。如果将所有逻辑揉在一个“大泥球”里,带来的后果是灾难性的:
- 变更影响面不可控:修改一个工具的实现可能波及整个工作流
- 难以独立扩展:客服业务暴涨时无法单独扩容对应模块
- 多厂商切换困难:某模型API升级时,需要到处修改调用代码
- 可观测性缺失:无法定位问题是出在模型推理、工具执行还是数据检索
分层架构的核心思想是通过职责分离来降低耦合度。借鉴OpenClaw等成熟AI网关的设计经验,我们将系统划分为三个核心层次:网关层、调度层、执行层。
┌─────────────────────────────────────────────────────────────┐
│ 接入层 │
│ (Webhook / WebSocket / REST API) │
└─────────────────────┬───────────────────────────────────────┘
▼
┌─────────────────────────────────────────────────────────────┐
│ 网关层 (Gateway) │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────┐ │
│ │ 路由与鉴权 │ │ 限流与熔断 │ │ 多模型路由与故障迁移 │ │
│ └─────────────┘ └─────────────┘ └─────────────────────┘ │
└─────────────────────┬───────────────────────────────────────┘
▼
┌─────────────────────────────────────────────────────────────┐
│ 调度层 (Orchestrator) │
│ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────┐ │
│ │ 任务拆解 │ │ 状态管理 │ │ 上下文编排与RAG检索 │ │
│ └─────────────┘ └─────────────┘ └─────────────────────┘ │
└─────────────────────┬───────────────────────────────────────┘
▼
┌─────────────────────────────────────────────────────────────┐
│ 执行层 (Executor) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────────┐ │
│ │ 客服Agent │ │ 广告Agent │ │ 供应链Agent│ │ 知识库检索 │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────────┘ │
└─────────────────────────────────────────────────────────────┘
二、网关层:统一流量入口与多模型治理
网关层是系统的“门面”,负责处理所有外部请求的接入、鉴权、限流和路由。
2.1 统一接入与协议适配
跨境电商场景下,消息来源多样:用户可能通过WhatsApp、Shopify Inbox、邮件或Web端发起咨询。网关层需要将这些不同协议的消息标准化为统一的内部事件格式。
# gateway/models.py
from dataclasses import dataclass
from enum import Enum
from typing import Optional, Dict, Any
from datetime import datetime
class ChannelType(Enum):
WHATSAPP = "whatsapp"
SHOPIFY = "shopify"
WEB = "web"
EMAIL = "email"
class MessagePriority(Enum):
P0 = "critical" # 退款、欺诈等紧急问题
P1 = "high" # 订单查询
P2 = "normal" # 常规咨询
@dataclass
class NormalizedMessage:
"""标准化后的消息对象"""
message_id: str
channel: ChannelType
user_id: str
session_id: str # 用于关联同一用户的连续对话
text: str
priority: MessagePriority
timestamp: datetime
metadata: Dict[str, Any] # 保留渠道特有字段(如WhatsApp的phone_number_id)
2.2 多模型API网关设计
跨境业务中,不同场景可能需要调用不同模型:客服对话用Claude(长上下文优势),广告文案生成用GPT-4(创意能力强),供应链查询用DeepSeek(成本敏感)。网关层需要统一管理这些模型API的调用、路由和故障迁移。
# gateway/model_router.py
from abc import ABC, abstractmethod
from typing import Dict, List, Optional
import httpx
import asyncio
from enum import Enum
class ModelProvider(Enum):
OPENAI = "openai"
ANTHROPIC = "anthropic"
DEEPSEEK = "deepseek"
GEMINI = "gemini"
class ModelCapability(Enum):
CHAT = "chat"
FUNCTION_CALLING = "function_calling"
LONG_CONTEXT = "long_context"
REASONING = "reasoning" # 推理能力,如o1系列
@dataclass
class ModelEndpoint:
provider: ModelProvider
model_name: str
base_url: str
api_key: str
capabilities: List[ModelCapability]
rate_limit: int # 每分钟请求数
weight: int = 10 # 路由权重
class ModelGateway:
"""
多模型API网关
支持按能力路由、权重负载均衡、故障自动迁移
"""
def __init__(self):
self.endpoints: Dict[str, List[ModelEndpoint]] = {}
self.failover_history: Dict[str, int] = {} # 记录故障迁移次数
self._init_endpoints()
def _init_endpoints(self):
# 实际应从配置中心加载
self.endpoints["chat"] = [
ModelEndpoint(
provider=ModelProvider.ANTHROPIC,
model_name="claude-sonnet-4-20250514",
base_url="https://api.anthropic.com/v1",
api_key=os.getenv("ANTHROPIC_API_KEY"),
capabilities=[ModelCapability.CHAT, ModelCapability.FUNCTION_CALLING],
rate_limit=100
),
ModelEndpoint(
provider=ModelProvider.DEEPSEEK,
model_name="deepseek-chat",
base_url="https://api.deepseek.com/v1",
api_key=os.getenv("DEEPSEEK_API_KEY"),
capabilities=[ModelCapability.CHAT, ModelCapability.FUNCTION_CALLING],
rate_limit=200,
weight=5 # 低成本模型权重更高
)
]
async def route_request(self, capability: ModelCapability,
messages: List[Dict],
tools: Optional[List[Dict]] = None) -> Dict:
"""
根据所需能力路由请求到合适的模型
支持自动重试和故障迁移
"""
candidates = self.endpoints.get("chat", [])
# 按能力过滤
candidates = [e for e in candidates if capability in e.capabilities]
# 按权重排序,实现加权轮询
candidates = sorted(candidates, key=lambda e: e.weight, reverse=True)
# 尝试调用,失败则迁移到下一个
for endpoint in candidates:
try:
return await self._call_model(endpoint, messages, tools)
except Exception as e:
print(f"Model {endpoint.model_name} failed: {e}")
# 记录故障,继续尝试下一个
key = f"{endpoint.provider.value}:{endpoint.model_name}"
self.failover_history[key] = self.failover_history.get(key, 0) + 1
continue
raise Exception("All model endpoints failed")
async def _call_model(self, endpoint: ModelEndpoint,
messages: List[Dict],
tools: Optional[List[Dict]] = None) -> Dict:
"""实际调用模型API"""
# 根据不同provider构造不同格式的请求
if endpoint.provider == ModelProvider.ANTHROPIC:
return await self._call_anthropic(endpoint, messages, tools)
elif endpoint.provider == ModelProvider.OPENAI:
return await self._call_openai(endpoint, messages, tools)
# ...
2.3 限流与配额管理
在多租户场景下,不同业务团队共享模型资源,需要基于Token消耗进行精细化限流,防止某个业务过度消耗影响其他业务。
# gateway/rate_limiter.py
import time
from collections import defaultdict
from typing import Dict, Tuple
class TokenBucketRateLimiter:
"""
基于Token消耗量的限流器
支持租户级别和模型级别的双重配额
"""
def __init__(self):
# tenant_quota: {tenant_id: {model: (tokens_limit, window_seconds)}}
self.tenant_quota: Dict[str, Dict[str, Tuple[int, int]]] = {}
# consumption记录: {tenant_id: {model: [(timestamp, tokens_used)]}}
self.consumption: Dict[str, Dict[str, List[Tuple[float, int]]]] = defaultdict(
lambda: defaultdict(list)
)
def check_and_consume(self, tenant_id: str, model: str, tokens: int) -> bool:
"""检查配额并消耗,返回是否允许"""
quota = self.tenant_quota.get(tenant_id, {}).get(model)
if not quota:
return True # 未配置配额则放行
limit, window = quota
now = time.time()
records = self.consumption[tenant_id][model]
# 清理过期的记录
records[:] = [(ts, t) for ts, t in records if now - ts < window]
# 计算当前窗口总消耗
total = sum(t for _, t in records)
if total + tokens > limit:
return False
records.append((now, tokens))
return True
三、调度层:任务拆解与状态管理
调度层是整个系统的“大脑”,负责将用户意图转化为可执行的任务序列,并管理整个执行过程的上下文状态。
3.1 Supervisor-Worker模式
借鉴分层Agent架构的设计思想,调度层采用Supervisor-Worker模式:一个Supervisor Agent负责任务拆解和路由决策,多个Worker Agent负责执行具体领域的任务。
# orchestrator/supervisor.py
from typing import List, Dict, Any, Optional
from langgraph.graph import StateGraph, END
from langgraph.checkpoint import MemorySaver
from pydantic import BaseModel
class TaskState(BaseModel):
"""任务状态模型"""
user_query: str
session_id: str
messages: List[Dict] = []
current_task: Optional[str] = None
sub_tasks: List[Dict] = []
results: Dict[str, Any] = {}
completed_steps: List[str] = []
error: Optional[str] = None
need_human_review: bool = False
class SupervisorAgent:
"""
调度中枢 - Supervisor Agent
负责任务拆解、路由分发和结果聚合
"""
def __init__(self, model_gateway, knowledge_retriever):
self.model_gateway = model_gateway
self.knowledge_retriever = knowledge_retriever
self._build_graph()
def _build_graph(self):
"""使用LangGraph构建工作流图"""
builder = StateGraph(TaskState)
# 添加节点
builder.add_node("analyze", self._analyze_intent)
builder.add_node("decompose", self._decompose_task)
builder.add_node("route", self._route_to_worker)
builder.add_node("aggregate", self._aggregate_results)
builder.add_node("human_review", self._human_review)
# 添加边
builder.set_entry_point("analyze")
builder.add_edge("analyze", "decompose")
builder.add_conditional_edges(
"decompose",
self._should_route,
{
"execute": "route",
"human": "human_review",
"end": END
}
)
builder.add_edge("route", "aggregate")
builder.add_edge("aggregate", END)
# 配置状态持久化
self.checkpointer = MemorySaver()
self.graph = builder.compile(checkpointer=self.checkpointer)
async def run(self, query: str, session_id: str) -> Dict:
"""执行任务"""
initial_state = TaskState(
user_query=query,
session_id=session_id,
messages=[{"role": "user", "content": query}]
)
config = {"configurable": {"thread_id": session_id}}
result = await self.graph.ainvoke(initial_state, config)
return result
async def _analyze_intent(self, state: TaskState) -> TaskState:
"""分析用户意图"""
prompt = f"""
分析以下用户问题的意图和紧急程度:
用户问题:{state.user_query}
请输出JSON格式:
{{
"intent": "订单查询|售后咨询|广告投放|供应链查询|其他",
"urgency": "P0|P1|P2",
"entities": {{"order_id": "...", "product_sku": "..."}}
}}
"""
response = await self.model_gateway.route_request(
capability=ModelCapability.CHAT,
messages=[{"role": "user", "content": prompt}]
)
# 解析响应并更新state
# ...
return state
async def _decompose_task(self, state: TaskState) -> TaskState:
"""将复杂任务拆解为子任务序列"""
# 使用Plan-and-Execute模式拆解任务
# ...
return state
async def _route_to_worker(self, state: TaskState) -> TaskState:
"""将子任务路由到对应的执行器"""
# 根据子任务类型调用对应的Worker
# ...
return state
3.2 状态持久化与时间旅行
跨境电商场景中,一个对话可能跨越数天(例如用户查询订单状态后,两天后追问物流进度)。调度层需要支持状态持久化,即每个会话的执行状态被完整保存,支持中断恢复和历史回溯。
LangGraph的状态持久化机制非常契合这一需求:每个步骤执行后,系统会自动保存当前图的状态快照(StateSnapshot),包括所有变量、消息历史、下一步要执行的节点等。通过thread_id(会话ID)区分不同会话,后续请求可以从此前中断的位置恢复执行。
# orchestrator/persistence.py
from langgraph.checkpoint import SqliteSaver
from langgraph.checkpoint.base import BaseCheckpointSaver
class StateManager:
"""
状态管理器
使用SQLite持久化会话状态,支持恢复和回溯
"""
def __init__(self, db_path: str = "checkpoints.db"):
self.checkpointer = SqliteSaver.from_conn_string(db_path)
def get_state(self, thread_id: str):
"""获取指定会话的当前状态"""
config = {"configurable": {"thread_id": thread_id}}
snapshot = self.checkpointer.get_tuple(config)
if snapshot:
return snapshot.checkpoint
return None
def get_history(self, thread_id: str, limit: int = 10):
"""获取会话的历史状态快照,用于调试和时间旅行"""
config = {"configurable": {"thread_id": thread_id}}
history = list(self.checkpointer.list(config, limit=limit))
return history
四、执行层:专业Agent与工具集成
执行层包含各类专业Worker Agent,每个Worker专注于特定领域,通过调用外部API或执行本地逻辑来完成具体任务。
4.1 工具注册与调用
每个Worker暴露一组工具(Tools),通过标准化的函数签名供上层调用。参考Function Calling的设计模式,工具以JSON Schema格式描述,便于大模型理解。
# executor/worker_base.py
from abc import ABC, abstractmethod
from typing import Dict, Any, List, Callable
import json
class Tool:
"""工具定义"""
def __init__(self, name: str, description: str, parameters: Dict, func: Callable):
self.name = name
self.description = description
self.parameters = parameters # JSON Schema格式
self.func = func
def to_openai_schema(self) -> Dict:
"""转换为OpenAI Function Calling格式"""
return {
"type": "function",
"function": {
"name": self.name,
"description": self.description,
"parameters": self.parameters
}
}
async def execute(self, **kwargs) -> Any:
"""执行工具"""
return await self.func(**kwargs)
class WorkerAgent(ABC):
"""Worker基类"""
def __init__(self, name: str):
self.name = name
self.tools: Dict[str, Tool] = {}
self._register_tools()
@abstractmethod
def _register_tools(self):
"""注册该Worker提供的工具"""
pass
def get_tools(self) -> List[Tool]:
"""获取所有工具"""
return list(self.tools.values())
async def execute_tool(self, tool_name: str, params: Dict) -> Any:
"""执行指定工具"""
if tool_name not in self.tools:
raise ValueError(f"Tool {tool_name} not found")
return await self.tools[tool_name].execute(**params)
4.2 跨境电商专用Worker示例
# executor/order_worker.py
from executor.worker_base import WorkerAgent, Tool
import httpx
class OrderWorker(WorkerAgent):
"""订单查询Worker - 对接ERP系统"""
def __init__(self):
super().__init__("order_worker")
def _register_tools(self):
self.tools["query_order"] = Tool(
name="query_order",
description="查询订单状态和详情,支持Amazon、Walmart、Shopify订单",
parameters={
"type": "object",
"properties": {
"order_id": {"type": "string", "description": "订单编号"},
"platform": {
"type": "string",
"enum": ["amazon", "walmart", "shopify"],
"description": "订单来源平台"
}
},
"required": ["order_id", "platform"]
},
func=self._query_order
)
self.tools["track_logistics"] = Tool(
name="track_logistics",
description="查询订单物流追踪信息",
parameters={
"type": "object",
"properties": {
"tracking_number": {"type": "string", "description": "物流单号"},
"carrier": {"type": "string", "enum": ["fedex", "ups", "dhl", "usps"]}
},
"required": ["tracking_number"]
},
func=self._track_logistics
)
async def _query_order(self, order_id: str, platform: str) -> Dict:
"""对接ERP系统查询订单"""
# 实际应通过统一的ERP适配器调用
# 这里仅做示例
async with httpx.AsyncClient() as client:
response = await client.get(
f"https://erp.company.com/api/orders/{order_id}",
params={"platform": platform},
headers={"Authorization": f"Bearer {os.getenv('ERP_API_KEY')}"}
)
return response.json()
async def _track_logistics(self, tracking_number: str, carrier: str) -> Dict:
"""对接物流API查询追踪信息"""
# 调用第三方物流API
# ...
pass
五、层间通信与解耦设计
三层的解耦依赖于明确定义的接口契约:
| 通信方向 | 数据格式 | 协议 |
|---|---|---|
| 接入层 → 网关层 | 标准化消息对象 | REST / WebSocket |
| 网关层 → 调度层 | (会话ID + 用户消息 + 元数据) | 内部gRPC / 消息队列 |
| 调度层 → 执行层 | (任务类型 + 参数 + 上下文) | 工具调用接口 |
这种设计带来的收益:
- 独立演进:各层可独立升级,只要接口不变即可
- 弹性扩缩:网关层可水平扩展处理高并发,执行层可按业务域独立扩容
- 故障隔离:某个Worker崩溃不影响网关和调度层
- 可观测性:每层独立记录日志,便于定位问题
六、总结
跨境AI智能体的分层架构设计,本质上是对复杂系统进行关注点分离的工程实践。网关层解决“怎么接进来”的问题(统一入口、模型路由、限流熔断),调度层解决“要做什么”的问题(任务拆解、状态管理、工作流编排),执行层解决“怎么做到”的问题(专业工具调用、系统对接)。
在实际落地中,三层的边界并非一成不变,需要根据业务规模和团队组织灵活调整。但核心理念始终不变:让每一层只做一件事,并把它做好。
更多推荐



所有评论(0)