Skip to content

阶段 4:多 Agent 协作 🔥 推荐 ​

一句话总结:一个 Agent 做不了的事,多个 Agent 分工协作就能做。

📊 学习进度

  • 状态:⬜ 未开始
  • 预计时长:4-5 小时
  • 已完成:0/3 个模块
  • 在整体流程中的位置:AI Agent 开发·第 4 阶段

📍 本章定位

  • 服务方案:方案 3(核心 80%)
  • 学习方式:🔥 推荐
  • 在流程中的作用:搭建多 Agent 协作系统,是方案 3 量化系统的核心架构
  • 核心知识点:Agent 分工、通信机制、协作模式
  • 预计时长:4-5 小时
  • 完成后能做什么:能设计和实现多 Agent 协作系统

1. 传统模式:痛点与瓶颈 ​

1.1 单 Agent 的局限性 ​

单 Agent 系统在处理简单任务时效率很高,但面对复杂、多领域、需要专业分工的任务时,存在根本性局限:

局限性说明量化数据
能力上限单一 Prompt 难以覆盖所有领域任务成功率下降 35%
上下文压力所有信息塞入一个上下文窗口Token 消耗增加 2.8x
错误累积单点故障导致整个任务失败错误率 18%
专业深度通用 Prompt 难以达到专业水平质量评分下降 40%

数据来源:LangGraph 2025 年多 Agent 基准测试

1.2 量化痛点数据 ​

据 CrewAI 2025 年基准测试:

指标单 Agent多 Agent改善幅度
复杂任务完成率45%87%+93%
专业任务质量3.2/54.6/5+44%
并行处理能力1 个任务5-10 个任务+500%
错误恢复能力单点故障自动降级+300%
系统可扩展性线性增长指数增长+200%

1.3 OPC 场景下的核心矛盾 ​

OPC 运营者需要同时处理多个专业领域的工作(数据分析、内容创作、客户服务、交易执行),单 Agent 无法同时精通所有领域。多 Agent 系统允许每个 Agent 专注一个领域,通过协作完成复杂任务。


2. OPC 模式:重新定义 ​

2.1 核心理念 ​

多 Agent 协作 = 专业分工 + 通信协调 + 共享状态。就像一个高效的团队,每个成员专注自己的领域,通过明确的协作机制完成复杂任务。

多 Agent 协作架构:

2.2 人机分工矩阵 ​

任务人类角色AI 角色协作方式
Agent 分工设计定义每个 Agent 的职责建议最佳分工人决策,AI 辅助
通信机制设计定义消息格式和流程实现通信逻辑人定义协议,AI 编码
协作模式选择选择层级/辩论/流水线实现对应模式人决策,AI 实现
冲突解决策略定义优先级和仲裁规则实现冲突检测人定义规则,AI 执行
性能监控设定 SLA 指标实时监控告警人审核,AI 执行

2.3 效率对比 ​

指标单 Agent多 Agent提效倍数
复杂任务完成率45%87%1.9x
专业任务质量3.2/54.6/51.4x
并行处理能力1 任务5-10 任务5-10x
系统可扩展性线性指数2x+

3. 实操案例 ​

3.1 场景描述 ​

场景:为 OPC 搭建一个"量化交易多 Agent 系统",包含以下 Agent:

  1. 数据 Agent — 采集链上和市场数据
  2. 分析 Agent — 技术面和基本面分析
  3. 策略 Agent — 生成交易信号
  4. 执行 Agent — 执行交易操作
  5. 风控 Agent — 监控风险并触发止损

技术栈:LangGraph + Claude API + Python

3.2 执行过程 ​

3.2.1 多 Agent 协作模式详解 ​

三种主流协作模式:

模式适用场景优点缺点
层级式任务可分解、需要统一调度职责清晰、易于管理主管成为瓶颈
辩论式需要多角度分析、决策减少偏见、提高质量耗时长、成本高
流水线式步骤固定、顺序执行简单直观、易于调试灵活性差

3.2.2 LangGraph 实现 ​

LangGraph 使用状态图(StateGraph)定义多 Agent 的协作流程:

Python 实现:

python
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode
import operator

# 定义状态
class TradingState(TypedDict):
    """交易系统状态"""
    market_data: dict  # 市场数据
    analysis: dict  # 分析结果
    signal: dict  # 交易信号
    risk_check: dict  # 风控检查
    execution_result: dict  # 执行结果
    messages: Annotated[list, operator.add]  # 消息历史

# 定义各 Agent 节点
def data_agent(state: TradingState) -> dict:
    """数据 Agent - 采集市场数据"""
    # 调用 LLM 分析数据需求
    market_data = {
        "BTC": {"price": 67500, "volume": 28.5, "change_24h": 2.3},
        "ETH": {"price": 3500, "volume": 15.2, "change_24h": 1.8}
    }
    return {
        "market_data": market_data,
        "messages": [f"数据 Agent: 已采集 {len(market_data)} 个币种数据"]
    }

def analysis_agent(state: TradingState) -> dict:
    """分析 Agent - 技术面分析"""
    market_data = state["market_data"]
    
    # 技术分析逻辑
    analysis = {}
    for symbol, data in market_data.items():
        analysis[symbol] = {
            "trend": "bullish" if data["change_24h"] > 0 else "bearish",
            "volume_signal": "high" if data["volume"] > 20 else "normal",
            "recommendation": "buy" if data["change_24h"] > 2 else "hold"
        }
    
    return {
        "analysis": analysis,
        "messages": [f"分析 Agent: 完成 {len(analysis)} 个币种分析"]
    }

def strategy_agent(state: TradingState) -> dict:
    """策略 Agent - 生成交易信号"""
    analysis = state["analysis"]
    
    signals = {}
    for symbol, info in analysis.items():
        if info["recommendation"] == "buy" and info["volume_signal"] == "high":
            signals[symbol] = {
                "action": "buy",
                "confidence": 0.85,
                "size": "small"  # 小仓位试探
            }
    
    return {
        "signal": signals,
        "messages": [f"策略 Agent: 生成 {len(signals)} 个交易信号"]
    }

def risk_agent(state: TradingState) -> dict:
    """风控 Agent - 风险检查"""
    signals = state["signal"]
    
    risk_check = {}
    for symbol, signal in signals.items():
        # 风控规则
        risk_check[symbol] = {
            "approved": signal["confidence"] > 0.7,
            "max_position": 0.1,  # 最大仓位 10%
            "stop_loss": 0.05  # 止损 5%
        }
    
    return {
        "risk_check": risk_check,
        "messages": [f"风控 Agent: 完成 {len(risk_check)} 个信号风控检查"]
    }

def execution_agent(state: TradingState) -> dict:
    """执行 Agent - 执行交易"""
    signals = state["signal"]
    risk_check = state["risk_check"]
    
    results = {}
    for symbol, signal in signals.items():
        if risk_check[symbol]["approved"]:
            # 实际项目中调用交易所 API
            results[symbol] = {
                "status": "executed",
                "action": signal["action"],
                "price": state["market_data"][symbol]["price"]
            }
    
    return {
        "execution_result": results,
        "messages": [f"执行 Agent: 执行 {len(results)} 笔交易"]
    }

# 构建状态图
workflow = StateGraph(TradingState)

# 添加节点
workflow.add_node("data", data_agent)
workflow.add_node("analysis", analysis_agent)
workflow.add_node("strategy", strategy_agent)
workflow.add_node("risk", risk_agent)
workflow.add_node("execution", execution_agent)

# 定义边(执行流程)
workflow.set_entry_point("data")
workflow.add_edge("data", "analysis")
workflow.add_edge("analysis", "strategy")
workflow.add_edge("strategy", "risk")

# 条件边:根据风控结果决定是否执行
def should_execute(state: TradingState) -> str:
    """检查是否有通过风控的信号"""
    approved = any(
        check["approved"] for check in state["risk_check"].values()
    )
    return "execution" if approved else END

workflow.add_conditional_edges("risk", should_execute)
workflow.add_edge("execution", END)

# 编译并运行
app = workflow.compile()

# 执行交易流程
result = app.invoke({
    "market_data": {},
    "analysis": {},
    "signal": {},
    "risk_check": {},
    "execution_result": {},
    "messages": []
})

print("执行结果:", result["execution_result"])
print("消息历史:", result["messages"])

TypeScript 实现(使用 LangGraph.js):

typescript
import { StateGraph, END } from "@langchain/langgraph";
import { Annotation } from "@langchain/langgraph";

// 定义状态
const TradingState = Annotation.Root({
  marketData: Annotation<Record<string, unknown>>,
  analysis: Annotation<Record<string, unknown>>,
  signal: Annotation<Record<string, unknown>>,
  riskCheck: Annotation<Record<string, unknown>>,
  executionResult: Annotation<Record<string, unknown>>,
  messages: Annotation<string[]>({
    reducer: (prev, next) => [...prev, ...next],
    default: () => [],
  }),
});

// 定义 Agent 节点
async function dataAgent(
  state: typeof TradingState.State
): Promise<Partial<typeof TradingState.State>> {
  const marketData = {
    BTC: { price: 67500, volume: 28.5, change24h: 2.3 },
    ETH: { price: 3500, volume: 15.2, change24h: 1.8 },
  };
  return {
    marketData,
    messages: [`数据 Agent: 已采集 ${Object.keys(marketData).length} 个币种数据`],
  };
}

async function analysisAgent(
  state: typeof TradingState.State
): Promise<Partial<typeof TradingState.State>> {
  const analysis: Record<string, unknown> = {};
  for (const [symbol, data] of Object.entries(state.marketData)) {
    const d = data as { change24h: number; volume: number };
    analysis[symbol] = {
      trend: d.change24h > 0 ? "bullish" : "bearish",
      recommendation: d.change24h > 2 ? "buy" : "hold",
    };
  }
  return {
    analysis,
    messages: [`分析 Agent: 完成分析`],
  };
}

// ... 其他 Agent 类似

// 构建状态图
const workflow = new StateGraph(TradingState)
  .addNode("data", dataAgent)
  .addNode("analysis", analysisAgent)
  // .addNode("strategy", strategyAgent)
  // .addNode("risk", riskAgent)
  // .addNode("execution", executionAgent)
  .addEdge("__start__", "data")
  .addEdge("data", "analysis")
  .addEdge("analysis", END);

const app = workflow.compile();
const result = await app.invoke({
  marketData: {},
  analysis: {},
  signal: {},
  riskCheck: {},
  executionResult: {},
  messages: [],
});

3.2.3 CrewAI 实现 ​

CrewAI 使用"角色扮演"模式定义多 Agent 协作:

python
from crewai import Agent, Task, Crew

# 定义 Agent(角色)
data_analyst = Agent(
    role="数据分析师",
    goal="采集和分析加密货币市场数据",
    backstory="你是一位经验丰富的加密货币数据分析师,擅长从链上数据中发现趋势。",
    tools=[search_tool, api_tool],
    llm="claude-sonnet-4-20250514"
)

strategy_expert = Agent(
    role="策略专家",
    goal="基于数据分析生成交易策略",
    backstory="你是一位量化交易策略专家,擅长将数据分析转化为可执行的交易信号。",
    tools=[calculator_tool],
    llm="claude-sonnet-4-20250514"
)

risk_manager = Agent(
    role="风控经理",
    goal="评估交易策略的风险并提供风控建议",
    backstory="你是一位严谨的风控专家,确保每笔交易都在可接受的风险范围内。",
    tools=[],
    llm="claude-sonnet-4-20250514"
)

# 定义任务
data_task = Task(
    description="分析 BTC 和 ETH 的最新市场数据,包括价格、成交量、链上活跃度。",
    expected_output="包含价格趋势、成交量分析、链上指标的结构化报告",
    agent=data_analyst
)

strategy_task = Task(
    description="基于数据分析结果,生成今日的交易策略建议。",
    expected_output="包含买入/卖出信号、仓位建议、止损位的策略报告",
    agent=strategy_expert,
    context=[data_task]  # 依赖数据任务的结果
)

risk_task = Task(
    description="评估交易策略的风险,检查是否符合风控规则。",
    expected_output="风险评估报告,包含风险等级、建议调整、最终批准/拒绝",
    agent=risk_manager,
    context=[strategy_task]
)

# 组建团队
crew = Crew(
    agents=[data_analyst, strategy_expert, risk_manager],
    tasks=[data_task, strategy_task, risk_task],
    verbose=True
)

# 执行任务
result = crew.kickoff()
print(result)

3.2.4 AutoGen 实现(微软) ​

AutoGen 使用"对话式"模式,Agent 之间通过消息对话协作:

python
from autogen import AssistantAgent, UserProxyAgent

# 配置 LLM
config = {
    "model": "claude-sonnet-4-20250514",
    "api_key": "your-api-key",
    "api_type": "anthropic"
}

# 定义分析 Agent
analyst = AssistantAgent(
    name="分析师",
    system_message="""你是一位加密货币市场分析师。
    分析数据时关注:价格趋势、成交量变化、链上活跃度。
    输出结构化分析报告。""",
    llm_config=config
)

# 定义策略 Agent
strategist = AssistantAgent(
    name="策略师",
    system_message="""你是一位量化交易策略专家。
    基于分析师的报告生成交易信号。
    每个信号包含:动作、置信度、仓位建议、止损位。""",
    llm_config=config
)

# 定义风控 Agent
risk_manager = AssistantAgent(
    name="风控经理",
    system_message="""你是一位风控专家。
    审查策略师的交易信号,确保:
    1. 单笔仓位不超过总资金 10%
    2. 止损设置合理(不超过 5%)
    3. 每日最大亏损不超过 2%
    输出批准/拒绝及理由。""",
    llm_config=config
)

# 定义用户代理(触发执行)
user_proxy = UserProxyAgent(
    name="用户",
    human_input_mode="NEVER",  # 无人值守模式
    max_consecutive_auto_reply=5,
    code_execution_config=False
)

# 启动多 Agent 对话
user_proxy.initiate_chat(
    analyst,
    message="分析 BTC 和 ETH 最近 24 小时的市场数据,生成今日交易建议。"
)

# AutoGen 会自动在 Agent 之间传递消息,直到达成共识

AutoGen vs LangGraph vs CrewAI 对比:

维度AutoGenLangGraphCrewAI
协作模式对话式状态图角色扮演
学习曲线中高低
灵活性高最高中
生产就绪中高中
适用场景复杂推理复杂工作流快速原型

3.2.5 Peer Review 模式 ​

Peer Review 模式让多个 Agent 互相审查对方的输出,提高决策质量。据 Microsoft Research 2025 年测试,Peer Review 模式在复杂推理任务上的准确率比单 Agent 高 28%。

python
from typing import Dict, List

class PeerReviewSystem:
    """Agent Peer Review 系统"""

    def __init__(self, agents: Dict[str, callable]):
        self.agents = agents

    async def review_cycle(
        self,
        task: str,
        max_rounds: int = 3,
        consensus_threshold: float = 0.8
    ) -> Dict:
        """执行 Peer Review 循环"""
        proposals = {}

        # 第1轮:每个 Agent 独立生成方案
        for name, agent in self.agents.items():
            proposals[name] = await agent(task)

        # 后续轮次:互相审查
        for round_num in range(max_rounds - 1):
            reviews = {}
            for reviewer_name, reviewer in self.agents.items():
                for proposer_name, proposal in proposals.items():
                    if reviewer_name == proposer_name:
                        continue
                    review = await reviewer(
                        f"请审查以下方案,指出优点和问题:\n{proposal}"
                    )
                    reviews[(reviewer_name, proposer_name)] = review

            # 检查共识度
            consensus = self._calculate_consensus(proposals, reviews)
            if consensus >= consensus_threshold:
                break

            # 根据反馈修订方案
            for name, agent in self.agents.items():
                feedback = [
                    r for (reviewer, proposer), r in reviews.items()
                    if proposer == name
                ]
                proposals[name] = await agent(
                    f"原始方案:{proposals[name]}\n\n审查反馈:{feedback}\n\n请修订方案。"
                )

        return {
            "proposals": proposals,
            "consensus": consensus,
            "rounds": round_num + 1
        }

    def _calculate_consensus(self, proposals: Dict, reviews: Dict) -> float:
        """计算共识度(简化版:基于正面评价比例)"""
        positive = sum(1 for r in reviews.values() if "同意" in r or "批准" in r)
        total = len(reviews)
        return positive / total if total > 0 else 0

3.2.6 Agent 通信机制 ​

多 Agent 之间的通信机制是协作的核心:

通信方式说明适用场景
共享状态所有 Agent 读写同一个状态对象LangGraph 默认方式
消息传递Agent 通过消息队列通信松耦合系统
事件驱动Agent 发布/订阅事件异步协作
直接调用Agent 直接调用其他 Agent紧耦合系统

共享状态示例:

python
from dataclasses import dataclass, field
from typing import Dict, Any
import threading

@dataclass
class SharedState:
    """多 Agent 共享状态"""
    data: Dict[str, Any] = field(default_factory=dict)
    lock: threading.Lock = field(default_factory=threading.Lock)
    
    def update(self, agent_id: str, key: str, value: Any):
        """Agent 更新共享状态"""
        with self.lock:
            if agent_id not in self.data:
                self.data[agent_id] = {}
            self.data[agent_id][key] = value
    
    def read(self, agent_id: str, key: str = None) -> Any:
        """Agent 读取共享状态"""
        with self.lock:
            if key:
                return self.data.get(agent_id, {}).get(key)
            return self.data.get(agent_id, {})
    
    def read_all(self) -> Dict[str, Any]:
        """读取所有 Agent 的状态"""
        with self.lock:
            return dict(self.data)

# 使用示例
state = SharedState()

# 数据 Agent 写入数据
state.update("data_agent", "market_data", {"BTC": 67500})

# 分析 Agent 读取数据
market_data = state.read("data_agent", "market_data")

3.2.7 冲突解决机制 ​

多 Agent 系统中,Agent 之间可能产生分歧。有效的冲突解决机制是系统可靠运行的关键。

四种冲突解决策略:

策略原理适用场景实现复杂度
主管裁决由主管 Agent 做最终决策层级式架构低
投票表决多数 Agent 同意即通过民主式决策中
置信度加权置信度高的 Agent 意见权重更大专业分工场景中
辩论收敛多轮辩论直到达成共识复杂推理任务高

置信度加权冲突解决实现:

python
from typing import Dict, List, Tuple
from dataclasses import dataclass

@dataclass
class AgentDecision:
    agent_id: str
    decision: str
    confidence: float  # 0-1
    reasoning: str

class ConflictResolver:
    """冲突解决器 - 置信度加权"""

    def resolve(self, decisions: List[AgentDecision]) -> Dict:
        """解决多个 Agent 之间的冲突"""
        if not decisions:
            return {"decision": None, "method": "no_decisions"}

        # 1. 按决策分组
        groups: Dict[str, List[AgentDecision]] = {}
        for d in decisions:
            if d.decision not in groups:
                groups[d.decision] = []
            groups[d.decision].append(d)

        # 2. 如果只有一个决策,直接返回
        if len(groups) == 1:
            return {
                "decision": decisions[0].decision,
                "method": "unanimous",
                "confidence": 1.0
            }

        # 3. 计算每个决策的加权得分
        scores: Dict[str, float] = {}
        for decision, agents in groups.items():
            # 加权得分 = 平均置信度 * 支持者比例
            avg_confidence = sum(a.confidence for a in agents) / len(agents)
            support_ratio = len(agents) / len(decisions)
            scores[decision] = avg_confidence * support_ratio

        # 4. 选择得分最高的决策
        best = max(scores.items(), key=lambda x: x[1])

        return {
            "decision": best[0],
            "method": "confidence_weighted",
            "confidence": best[1],
            "breakdown": {
                d: {"score": s, "supporters": len(groups[d])}
                for d, s in scores.items()
            }
        }

# 使用示例
resolver = ConflictResolver()
decisions = [
    AgentDecision("analyst", "buy", 0.85, "技术指标显示上涨趋势"),
    AgentDecision("strategist", "buy", 0.70, "基本面支持看涨"),
    AgentDecision("risk_manager", "hold", 0.60, "市场波动性较高"),
]
result = resolver.resolve(decisions)
# 结果: {"decision": "buy", "method": "confidence_weighted", "confidence": 0.57}

3.3 前后对比 ​

维度单 Agent多 Agent改善
复杂任务完成率45%87%+93%
专业任务质量3.2/54.6/5+44%
并行处理能力1 任务5-10 任务+500%
系统可扩展性线性指数+200%
错误恢复单点故障自动降级+300%

4. 趋势预判(未来 1-3 年) ​

4.1 技术演进方向 ​

技术方向当前状态1 年后3 年后
多 Agent 框架3-4 家主导整合为 2-3 家标准化
Agent 通信协议各自为政MCP 扩展统一标准
Agent 安全基础完善强制审计
Agent 市场萌芽增长繁荣

4.2 角色变化趋势 ​

角色当前1 年后3 年后
多 Agent 架构师稀缺主流基础能力
Agent 协作设计师几乎不存在新兴岗位标准配置
Agent 安全审计缺失初步建立强制要求

4.3 OPC 需要提前准备的能力 ​

  1. 状态图设计:理解 LangGraph 的状态机模型
  2. Agent 分工:如何将复杂任务分解为多个专业 Agent
  3. 通信协议:设计 Agent 间的消息格式和交互流程
  4. 冲突解决:多 Agent 意见不一致时的仲裁机制

5. 核心洞察 ​

🔑 关键洞察

多 Agent 系统的核心价值不是"更多 Agent",而是"专业分工"。一个由 5 个专业 Agent 组成的系统,在各自领域的表现可以超越一个通用 Agent 2-3 倍。关键是让每个 Agent 专注自己擅长的领域,而不是让所有 Agent 都做同样的事。

⚠️ 协作成本警告

多 Agent 系统的通信成本不可忽视。每增加一个 Agent,通信开销增加约 30%。建议从 2-3 个 Agent 开始,验证协作效果后再扩展。一个 10 个 Agent 的系统,如果设计不当,可能比单 Agent 更慢更贵。


6. 参考与延伸 ​

[1] LangGraph. "Documentation" — 多 Agent 编排框架(2025)

[2] CrewAI. "Documentation" — 角色扮演多 Agent 框架(2025)

[3] AutoGen. "Microsoft Research" — 对话式多 Agent 框架(2025)

[4] OpenAI. "Agents SDK" — OpenAI 多 Agent 框架(2025)

[5] Anthropic. "Multi-Agent Patterns" — 多 Agent 设计模式(2025)

[6] Wu et al. "AutoGen: Enabling Next-Gen LLM Applications" — AutoGen 论文(2023)

[7] Microsoft Research. "Multi-Agent Debate" — 多 Agent 辩论论文(2023)

[8] Gartner. "Top Strategic Technology Trends 2025" — Agent 技术趋势预测(2025)

[9] LangGraph. "Multi-Agent Benchmark" — 多 Agent 基准测试(2025)

[10] Anthropic. "Building Effective Agents" — Agent 构建指南(2025)

[11] CrewAI. "Performance Optimization" — 性能优化文档(2025)

[12] OpenAI. "Multi-Agent Patterns" — 多 Agent 设计模式(2025)


常见问题 ​

问题原因解决方案
Agent 冲突职责不清明确分工边界,定义优先级
通信开销大消息太多消息压缩 + 过滤
死锁循环依赖超时机制 + 依赖分析
成本失控Agent 数量过多从 2-3 个开始,逐步扩展
调试困难分布式执行使用 LangSmith 追踪

多 Agent 性能优化 ​

Agent 数量与性能关系 ​

多 Agent 系统的性能并非随 Agent 数量线性增长。过多的 Agent 会导致通信开销急剧增加,反而降低整体效率。

Agent 数量与性能关系:

Agent 数量任务完成率平均延迟通信开销推荐度
145%基准无简单任务
2-378%+20%低推荐
4-587%+45%中复杂任务
6-889%+80%高特定场景
10+85%+150%极高不推荐

数据来源:LangGraph 2025 年多 Agent 基准测试

关键发现:

  • 2-3 个 Agent 是性价比最高的配置
  • 超过 5 个 Agent 后,性能提升趋于平缓
  • 超过 8 个 Agent 后,通信开销导致性能下降

通信优化策略 ​

优化策略说明效果
消息压缩压缩 Agent 间传递的消息通信量减少 40-60%
消息过滤只传递必要信息通信量减少 30-50%
异步通信非阻塞消息传递延迟降低 50-70%
批量传递合并多条消息一次传递通信次数减少 60-80%

消息压缩实现:

python
import zlib
import json
from typing import Any

class MessageCompressor:
    """Agent 间消息压缩器"""

    def compress(self, message: dict) -> bytes:
        """压缩消息"""
        # 1. 序列化为 JSON
        json_str = json.dumps(message, ensure_ascii=False)

        # 2. 压缩
        compressed = zlib.compress(json_str.encode('utf-8'), level=6)

        return compressed

    def decompress(self, compressed: bytes) -> dict:
        """解压消息"""
        # 1. 解压
        json_bytes = zlib.decompress(compressed)

        # 2. 反序列化
        return json.loads(json_bytes.decode('utf-8'))

    def compress_selective(self, message: dict, keep_keys: list) -> dict:
        """选择性压缩 - 只保留必要字段"""
        return {k: v for k, v in message.items() if k in keep_keys}

# 使用示例
compressor = MessageCompressor()

# 原始消息
message = {
    "agent": "data_agent",
    "type": "market_data",
    "data": {"BTC": 67500, "ETH": 3500},
    "metadata": {"timestamp": "2025-06-15", "source": "coinbase"}
}

# 压缩
compressed = compressor.compress(message)
print(f"原始大小: {len(json.dumps(message))} bytes")
print(f"压缩后大小: {len(compressed)} bytes")
print(f"压缩率: {len(compressed) / len(json.dumps(message)) * 100:.1f}%")

负载均衡策略 ​

当多个 Agent 可以处理同类任务时,负载均衡可以提高系统吞吐量和可靠性。

三种负载均衡策略:

策略原理优点缺点适用场景
轮询依次分配任务实现简单不考虑负载任务均匀
最少连接分配给最空闲的 Agent负载均衡需要状态跟踪任务不均匀
加权根据 Agent 能力分配能力匹配配置复杂异构 Agent

最少连接负载均衡实现:

python
from typing import Dict, List
from dataclasses import dataclass

@dataclass
class AgentStatus:
    agent_id: str
    active_tasks: int
    max_tasks: int
    success_rate: float

class LoadBalancer:
    """Agent 负载均衡器"""

    def __init__(self, agents: List[AgentStatus]):
        self.agents = {a.agent_id: a for a in agents}

    def select_agent(self, task_type: str = None) -> str:
        """选择最合适的 Agent"""
        # 过滤可用 Agent
        available = [
            a for a in self.agents.values()
            if a.active_tasks < a.max_tasks
        ]

        if not available:
            raise Exception("没有可用的 Agent")

        # 选择活跃任务最少的 Agent
        selected = min(available, key=lambda a: a.active_tasks)
        return selected.agent_id

    def update_status(self, agent_id: str, active_tasks: int):
        """更新 Agent 状态"""
        if agent_id in self.agents:
            self.agents[agent_id].active_tasks = active_tasks

    def get_stats(self) -> dict:
        """获取负载统计"""
        total_active = sum(a.active_tasks for a in self.agents.values())
        total_capacity = sum(a.max_tasks for a in self.agents.values())

        return {
            "agents": len(self.agents),
            "active_tasks": total_active,
            "total_capacity": total_capacity,
            "utilization": total_active / total_capacity * 100 if total_capacity > 0 else 0
        }

# 使用示例
balancer = LoadBalancer([
    AgentStatus("agent_1", active_tasks=2, max_tasks=5, success_rate=0.95),
    AgentStatus("agent_2", active_tasks=1, max_tasks=5, success_rate=0.98),
    AgentStatus("agent_3", active_tasks=3, max_tasks=5, success_rate=0.92),
])

# 选择 Agent
selected = balancer.select_agent()
print(f"选择 Agent: {selected}")

实操案例:量化交易多 Agent 系统优化 ​

场景:一个包含 5 个 Agent 的量化交易系统,需要优化以降低成本和提高响应速度。

优化前问题:

  • 月度 API 成本 $800
  • 平均决策延迟 45 秒
  • 任务成功率 82%

优化方案:

优化措施实施内容效果
Agent 精简将 5 个 Agent 合并为 3 个通信开销减少 40%
模型路由简单任务用 HaikuAPI 成本减少 60%
消息压缩压缩 Agent 间消息通信量减少 50%
并行执行数据采集并行化延迟减少 35%
缓存复用缓存市场数据查询API 调用减少 30%

优化后效果:

  • 月度 API 成本 $280(减少 65%)
  • 平均决策延迟 22 秒(减少 51%)
  • 任务成功率 91%(提升 11%)

下一步 ​

完成多 Agent 协作后,进入 阶段 5:部署运行 — 学习如何将多 Agent 系统部署到生产环境。

OPC 超级个体实战指南