一、为什么需要分层解耦

在跨境电商场景中,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 / 消息队列
调度层 → 执行层(任务类型 + 参数 + 上下文)工具调用接口

这种设计带来的收益:

  1. 独立演进:各层可独立升级,只要接口不变即可
  2. 弹性扩缩:网关层可水平扩展处理高并发,执行层可按业务域独立扩容
  3. 故障隔离:某个Worker崩溃不影响网关和调度层
  4. 可观测性:每层独立记录日志,便于定位问题

六、总结

跨境AI智能体的分层架构设计,本质上是对复杂系统进行关注点分离的工程实践。网关层解决“怎么接进来”的问题(统一入口、模型路由、限流熔断),调度层解决“要做什么”的问题(任务拆解、状态管理、工作流编排),执行层解决“怎么做到”的问题(专业工具调用、系统对接)。

在实际落地中,三层的边界并非一成不变,需要根据业务规模和团队组织灵活调整。但核心理念始终不变:让每一层只做一件事,并把它做好

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐