Python實(shí)現(xiàn)多智能體協(xié)作的實(shí)戰(zhàn)教程
一個(gè) Agent 解決不了的問(wèn)題,就組一個(gè)隊(duì)。本文帶你用 Python 從零搭建多智能體協(xié)作系統(tǒng),實(shí)現(xiàn)"AI 產(chǎn)品分析師 + 代碼工程師 + 測(cè)試專(zhuān)家 + 文檔撰寫(xiě)員 + 質(zhì)檢官"全自動(dòng)流水線(xiàn)。
為什么一個(gè) Agent 不夠用?
用過(guò) GPT 的人都有過(guò)這種經(jīng)歷:
你:幫我完整分析一個(gè)競(jìng)品,寫(xiě)一份技術(shù)方案,再給出完整代碼實(shí)現(xiàn),附上測(cè)試用例,最后寫(xiě)成報(bào)告……
GPT:好的?。ńo了一份殘缺不全、越寫(xiě)越錯(cuò)的東西)
單個(gè) Agent 的問(wèn)題在于:
- 上下文有限:長(zhǎng)任務(wù)超出 context window 就開(kāi)始"失憶"
- 角色模糊:既要分析又要寫(xiě)代碼,什么都做但什么都不精
- 無(wú)法并行:順序執(zhí)行,慢
多智能體協(xié)作的思路:像管理一個(gè)團(tuán)隊(duì),每個(gè) Agent 只做一件專(zhuān)業(yè)的事,通過(guò)協(xié)作完成復(fù)雜任務(wù)。
一、多智能體架構(gòu)設(shè)計(jì)
本文搭建的多智能體系統(tǒng)結(jié)構(gòu)如下:
用戶(hù)輸入任務(wù)
↓
Boss Agent(任務(wù)規(guī)劃 + 調(diào)度)
↓
┌──────────────────────────────────┐
│ 任務(wù)隊(duì)列(Task Queue) │
└──────────────────────────────────┘
↓ ↓ ↓
研究員 Agent 工程師 Agent 質(zhì)檢 Agent
(搜索/分析) (寫(xiě)代碼/方案) (測(cè)試/評(píng)審)
↓ ↓
文檔 Agent 測(cè)試 Agent
(寫(xiě)報(bào)告) (單元測(cè)試)
↓
最終輸出(報(bào)告 / 代碼 / 測(cè)試)
五位團(tuán)隊(duì)成員:
| Agent 角色 | 職責(zé) | 核心能力 |
|---|---|---|
| Boss Agent | 任務(wù)拆解 + 調(diào)度 + 匯總 | 規(guī)劃、協(xié)調(diào)、決策 |
| 研究員 Agent | 信息搜集 + 競(jìng)品分析 | 搜索、總結(jié)、報(bào)告 |
| 工程師 Agent | 代碼實(shí)現(xiàn) + 方案設(shè)計(jì) | 編程、架構(gòu)、調(diào)試 |
| 測(cè)試 Agent | 單元測(cè)試 + Bug 檢測(cè) | 測(cè)試、驗(yàn)證、覆蓋率 |
| 文檔 Agent | 技術(shù)文檔 + README | 寫(xiě)作、整理、規(guī)范 |
二、Agent 基類(lèi)封裝
# agents/base.py
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Any, Optional
from enum import Enum
import asyncio
import time
import uuid
import json
class TaskStatus(Enum):
PENDING = "pending"
RUNNING = "running"
DONE = "done"
FAILED = "failed"
CANCELLED = "cancelled"
@dataclass
class Task:
"""一個(gè)可執(zhí)行的任務(wù)單元"""
id: str = field(default_factory=lambda: str(uuid.uuid4())[:8])
name: str = ""
description: str = ""
assignee: str = "" # 分配給哪個(gè) Agent
depends_on: list[str] = field(default_factory=list) # 依賴(lài)的任務(wù) ID
status: TaskStatus = TaskStatus.PENDING
result: Any = None
error: str = ""
created_at: float = field(default_factory=time.time)
started_at: float = 0.0
finished_at: float = 0.0
retries: int = 0
max_retries: int = 3
@property
def elapsed_ms(self) -> float:
if self.finished_at and self.started_at:
return (self.finished_at - self.started_at) * 1000
return 0.0
def to_dict(self) -> dict:
return {
"id": self.id,
"name": self.name,
"assignee": self.assignee,
"status": self.status.value,
"elapsed_ms": round(self.elapsed_ms, 1),
"result_preview": str(self.result)[:200] if self.result else None
}
@dataclass
class AgentMemory:
"""Agent 的短期工作記憶"""
context: dict = field(default_factory=dict) # 鍵值對(duì)上下文
history: list[dict] = field(default_factory=list) # 對(duì)話(huà)歷史
def remember(self, key: str, value: Any):
self.context[key] = value
def recall(self, key: str, default=None) -> Any:
return self.context.get(key, default)
def add_message(self, role: str, content: str):
self.history.append({"role": role, "content": content})
def clear_history(self):
self.history.clear()
class BaseAgent(ABC):
"""
所有 Agent 的抽象基類(lèi)
每個(gè) Agent 有:角色、工具、記憶、LLM 后端
"""
def __init__(
self,
name: str,
role: str,
model: str = "gpt-4o",
api_key: str = None,
base_url: str = "https://api.openai.com/v1",
temperature: float = 0.7,
max_tokens: int = 4096,
verbose: bool = True
):
self.name = name
self.role = role
self.model = model
self.temperature = temperature
self.max_tokens = max_tokens
self.verbose = verbose
self.memory = AgentMemory()
# 初始化 LLM 客戶(hù)端
from openai import AsyncOpenAI
self.client = AsyncOpenAI(
api_key=api_key,
base_url=base_url
)
self._tools: dict[str, callable] = {}
self._register_tools()
def _register_tools(self):
"""子類(lèi)重寫(xiě),注冊(cè)本 Agent 可用的工具"""
pass
def register_tool(self, name: str, func: callable):
self._tools[name] = func
def log(self, msg: str):
if self.verbose:
print(f"[{self.name}] {msg}")
async def think(self, prompt: str, system_override: str = None) -> str:
"""調(diào)用 LLM 進(jìn)行思考,返回文本結(jié)果"""
system = system_override or f"""你是 {self.name},{self.role}。
請(qǐng)認(rèn)真完成分配給你的任務(wù),輸出要專(zhuān)業(yè)、詳細(xì)、有條理。"""
messages = [
{"role": "system", "content": system},
*self.memory.history[-6:], # 保留最近 6 條歷史(節(jié)省 token)
{"role": "user", "content": prompt}
]
response = await self.client.chat.completions.create(
model=self.model,
messages=messages,
temperature=self.temperature,
max_tokens=self.max_tokens
)
reply = response.choices[0].message.content
self.memory.add_message("assistant", reply)
return reply
@abstractmethod
async def execute(self, task: Task, shared_context: dict) -> Any:
"""執(zhí)行任務(wù),子類(lèi)必須實(shí)現(xiàn)"""
pass
三、五大專(zhuān)業(yè) Agent 實(shí)現(xiàn)
3.1 研究員 Agent(搜索 + 分析)
# agents/researcher.py
from .base import BaseAgent, Task
import httpx
class ResearcherAgent(BaseAgent):
"""研究員:負(fù)責(zé)信息搜集、競(jìng)品分析、數(shù)據(jù)整理"""
def __init__(self, **kwargs):
super().__init__(
name="研究員小研",
role="資深產(chǎn)品研究員,擅長(zhǎng)競(jìng)品分析、市場(chǎng)調(diào)研和信息整理",
temperature=0.3, # 研究任務(wù)需要更準(zhǔn)確,降低隨機(jī)性
**kwargs
)
def _register_tools(self):
self.register_tool("web_search", self._web_search)
self.register_tool("summarize_url", self._summarize_url)
async def _web_search(self, query: str) -> str:
"""模擬網(wǎng)頁(yè)搜索(實(shí)際接入 Serper / Tavily / DuckDuckGo)"""
# 實(shí)際使用時(shí)替換為真實(shí) API
async with httpx.AsyncClient() as client:
# 示例:接入 Serper API
# response = await client.post(
# "https://google.serper.dev/search",
# headers={"X-API-KEY": SERPER_API_KEY},
# json={"q": query}
# )
# return str(response.json()["organic"][:3])
return f"搜索結(jié)果({query}):[模擬數(shù)據(jù),請(qǐng)?zhí)鎿Q為真實(shí)搜索 API]"
async def _summarize_url(self, url: str) -> str:
"""抓取并總結(jié)網(wǎng)頁(yè)內(nèi)容"""
async with httpx.AsyncClient(timeout=10) as client:
try:
response = await client.get(url, follow_redirects=True)
# 簡(jiǎn)化處理,實(shí)際應(yīng)該用 BeautifulSoup 解析
text = response.text[:3000]
return await self.think(f"請(qǐng)總結(jié)以下網(wǎng)頁(yè)內(nèi)容的核心信息:\n{text}")
except Exception as e:
return f"抓取失敗: {e}"
async def execute(self, task: Task, shared_context: dict) -> str:
self.log(f"開(kāi)始執(zhí)行研究任務(wù):{task.name}")
task.started_at = __import__("time").time()
# 1. 搜集信息
search_result = await self._web_search(task.description)
self.memory.remember("search_result", search_result)
# 2. 深度分析
analysis_prompt = f"""
你是專(zhuān)業(yè)研究員。請(qǐng)根據(jù)以下信息,對(duì)「{task.name}」進(jìn)行深度分析:
搜索結(jié)果:
{search_result}
已有上下文:
{json.dumps(shared_context.get("background", {}), ensure_ascii=False, indent=2)}
請(qǐng)輸出:
1. **核心發(fā)現(xiàn)**(3-5 條)
2. **數(shù)據(jù)支撐**(關(guān)鍵數(shù)字和事實(shí))
3. **競(jìng)品對(duì)比**(表格形式)
4. **關(guān)鍵洞察**(對(duì)后續(xù)工作的啟示)
"""
result = await self.think(analysis_prompt)
self.log(f"研究完成,輸出 {len(result)} 字")
return result
3.2 工程師 Agent(寫(xiě)代碼 + 方案)
# agents/engineer.py
from .base import BaseAgent, Task
import subprocess, tempfile, os
class EngineerAgent(BaseAgent):
"""工程師:負(fù)責(zé)技術(shù)方案設(shè)計(jì)和代碼實(shí)現(xiàn)"""
def __init__(self, **kwargs):
super().__init__(
name="工程師小碼",
role="資深 Python 全棧工程師,擅長(zhǎng)架構(gòu)設(shè)計(jì)、代碼實(shí)現(xiàn)和性能優(yōu)化",
temperature=0.2, # 代碼任務(wù)要更確定,溫度調(diào)低
**kwargs
)
def _register_tools(self):
self.register_tool("run_code", self._run_code)
self.register_tool("check_syntax", self._check_syntax)
async def _run_code(self, code: str, timeout: int = 10) -> str:
"""在沙箱中執(zhí)行 Python 代碼"""
with tempfile.NamedTemporaryFile(mode="w", suffix=".py", delete=False) as f:
f.write(code)
f.flush()
try:
result = subprocess.run(
["python", f.name],
capture_output=True,
text=True,
timeout=timeout
)
output = result.stdout or result.stderr
return f"執(zhí)行結(jié)果:\n{output[:2000]}"
except subprocess.TimeoutExpired:
return "?? 執(zhí)行超時(shí)"
except Exception as e:
return f"? 執(zhí)行失?。簕e}"
finally:
os.unlink(f.name)
async def _check_syntax(self, code: str) -> str:
"""檢查 Python 代碼語(yǔ)法"""
try:
import ast
ast.parse(code)
return "? 語(yǔ)法檢查通過(guò)"
except SyntaxError as e:
return f"? 語(yǔ)法錯(cuò)誤:{e}"
def _extract_code(self, text: str) -> str:
"""從 LLM 輸出中提取代碼塊"""
import re
pattern = r"```(?:python)?\n(.*?)```"
matches = re.findall(pattern, text, re.DOTALL)
return matches[0].strip() if matches else text
async def execute(self, task: Task, shared_context: dict) -> dict:
self.log(f"開(kāi)始開(kāi)發(fā):{task.name}")
task.started_at = __import__("time").time()
research = shared_context.get("research_result", "")
# 1. 方案設(shè)計(jì)
design_prompt = f"""
任務(wù):{task.description}
參考研究結(jié)論:
{research[:1500] if research else "無(wú)"}
請(qǐng)先輸出:
## 技術(shù)方案
1. 架構(gòu)選型(為什么這么選)
2. 核心模塊拆解
3. 關(guān)鍵技術(shù)點(diǎn)
4. 預(yù)期接口設(shè)計(jì)
然后輸出完整的 Python 實(shí)現(xiàn)代碼。
"""
design_output = await self.think(design_prompt)
# 2. 提取并驗(yàn)證代碼
code = self._extract_code(design_output)
syntax_check = await self._check_syntax(code)
# 3. 如果語(yǔ)法有問(wèn)題,自動(dòng)修復(fù)
if "?" in syntax_check:
fix_prompt = f"以下代碼有語(yǔ)法錯(cuò)誤,請(qǐng)修復(fù):\n{syntax_check}\n\n```python\n[code]\n```"
fixed = await self.think(fix_prompt)
code = self._extract_code(fixed)
# 4. 運(yùn)行測(cè)試
run_result = await self._run_code(code)
self.log(f"代碼執(zhí)行:{run_result[:100]}")
return {
"design": design_output,
"code": code,
"syntax_check": syntax_check,
"run_result": run_result
}
3.3 測(cè)試 Agent + 文檔 Agent
# agents/tester.py
class TesterAgent(BaseAgent):
"""測(cè)試專(zhuān)家:生成單元測(cè)試,檢測(cè) Bug"""
def __init__(self, **kwargs):
super().__init__(
name="測(cè)試專(zhuān)家小測(cè)",
role="資深 QA 工程師,擅長(zhǎng)編寫(xiě) pytest 單元測(cè)試和發(fā)現(xiàn)潛在 Bug",
temperature=0.2,
**kwargs
)
async def execute(self, task: Task, shared_context: dict) -> dict:
self.log("開(kāi)始生成測(cè)試...")
code = shared_context.get("engineer_result", {}).get("code", "")
test_prompt = f"""
請(qǐng)為以下 Python 代碼編寫(xiě)全面的 pytest 測(cè)試套件:
```python
[code]
要求:
- 覆蓋所有公開(kāi)函數(shù)/方法
- 包含:正常用例、邊界條件、異常用例
- 使用 pytest.mark.parametrize 參數(shù)化測(cè)試
- 添加 Mock(如果涉及外部依賴(lài))
- 目標(biāo)覆蓋率:> 80%
輸出完整的可運(yùn)行測(cè)試文件。
test_code = await self.think(test_prompt)
# 分析測(cè)試覆蓋情況
analysis_prompt = f"分析以上測(cè)試的覆蓋情況,列出可能遺漏的測(cè)試用例:"
coverage_analysis = await self.think(analysis_prompt)
return {
"test_code": test_code,
"coverage_analysis": coverage_analysis
}agents/documenter.py
class DocumenterAgent(BaseAgent):
'文檔撰寫(xiě)員:生成技術(shù)文檔和 README'
def __init__(self, **kwargs):
super().__init__(
name="文檔小寫(xiě)",
role="技術(shù)文檔專(zhuān)家,擅長(zhǎng)撰寫(xiě)清晰易讀的 README、API 文檔和技術(shù)博客",
temperature=0.6, # 文檔寫(xiě)作可以稍微有創(chuàng)意
**kwargs
)
async def execute(self, task: Task, shared_context: dict) -> str:
self.log("開(kāi)始撰寫(xiě)文檔...")
engineer_result = shared_context.get("engineer_result", {})
tester_result = shared_context.get("tester_result", {})
research_result = shared_context.get("research_result", "")
doc_prompt = f"""請(qǐng)為以下項(xiàng)目生成完整的 README.md:
【項(xiàng)目背景】
{research_result[:500] if research_result else task.description}
【技術(shù)方案】
{engineer_result.get(“design”, “”)[:800]}
【代碼實(shí)現(xiàn)】
{engineer_result.get("code", "")[:1000]}
【測(cè)試情況】
{tester_result.get(“coverage_analysis”, “”)[:300]}README 需包含:
- 項(xiàng)目簡(jiǎn)介(一句話(huà))
- 功能特性(帶 emoji 的列表)
- 快速開(kāi)始(安裝 + 使用示例)
- API 參考(關(guān)鍵函數(shù)說(shuō)明)
- 項(xiàng)目結(jié)構(gòu)
- 貢獻(xiàn)指南
- License
風(fēng)格:專(zhuān)業(yè)但不枯燥,對(duì)新手友好。
四、任務(wù)調(diào)度器:Boss Agent
agents/boss.py
import asyncio
from dataclasses import dataclass, field
from .base import BaseAgent, Task, TaskStatus
class BossAgent(BaseAgent):
"""
Boss Agent:整體任務(wù)規(guī)劃 + 調(diào)度 + 結(jié)果匯總
相當(dāng)于項(xiàng)目經(jīng)理,負(fù)責(zé)把大任務(wù)拆成小任務(wù),分配給合適的 Agent
"""
def __init__(self, team: dict[str, BaseAgent], **kwargs):
super().__init__(
name="Boss 總監(jiān)",
role="資深項(xiàng)目總監(jiān),擅長(zhǎng)任務(wù)規(guī)劃、團(tuán)隊(duì)協(xié)調(diào)和結(jié)果整合",
temperature=0.5,
**kwargs
)
self.team = team # {agent_name: agent_instance}
self.task_queue: list[Task] = [] # 任務(wù)隊(duì)列
self.shared_context: dict = {} # 團(tuán)隊(duì)共享上下文
async def plan(self, user_request: str) -> list[Task]:
"""根據(jù)用戶(hù)需求,拆解任務(wù)并分配給團(tuán)隊(duì)成員"""
team_info = "\n".join([
f"- {name}: {agent.role}"
for name, agent in self.team.items()
])
plan_prompt = f"""
用戶(hù)需求:{user_request}
你的團(tuán)隊(duì)成員:
{team_info}
請(qǐng)把這個(gè)需求拆解成 3-6 個(gè)可執(zhí)行的子任務(wù),輸出 JSON 格式:
[
{{
"name": "任務(wù)名稱(chēng)",
"description": "詳細(xì)描述(給 Agent 的明確指令)",
"assignee": "agent 名稱(chēng)(researcher/engineer/tester/documenter)",
"depends_on": [] // 依賴(lài)的其他任務(wù)名稱(chēng)
}}
]
注意:
1. 每個(gè)任務(wù)職責(zé)單一,可獨(dú)立執(zhí)行
2. 正確設(shè)置依賴(lài)關(guān)系(如寫(xiě)代碼必須等研究完成)
3. 任務(wù)描述要具體,讓 Agent 知道該做什么
"""
plan_output = await self.think(plan_prompt)
# 提取 JSON 任務(wù)列表
import re, json
match = re.search(r"\[.*?\]", plan_output, re.DOTALL)
if match:
try:
tasks_data = json.loads(match.group())
return [Task(**t) for t in tasks_data]
except:
pass
# fallback:默認(rèn)任務(wù)流
return self._default_task_flow(user_request)
def _default_task_flow(self, request: str) -> list[Task]:
return [
Task(name="需求調(diào)研", description=f"調(diào)研:{request}", assignee="researcher"),
Task(name="代碼實(shí)現(xiàn)", description=f"實(shí)現(xiàn):{request}", assignee="engineer", depends_on=["需求調(diào)研"]),
Task(name="單元測(cè)試", description="為以上代碼生成測(cè)試", assignee="tester", depends_on=["代碼實(shí)現(xiàn)"]),
Task(name="文檔生成", description="生成完整文檔", assignee="documenter", depends_on=["代碼實(shí)現(xiàn)", "單元測(cè)試"]),
]
async def _execute_task(self, task: Task) -> None:
"""執(zhí)行單個(gè)任務(wù)"""
agent_map = {
"researcher": self.team.get("researcher"),
"engineer": self.team.get("engineer"),
"tester": self.team.get("tester"),
"documenter": self.team.get("documenter"),
}
agent = agent_map.get(task.assignee)
if not agent:
task.status = TaskStatus.FAILED
task.error = f"找不到 Agent:{task.assignee}"
return
task.status = TaskStatus.RUNNING
self.log(f"? 分派任務(wù)「{task.name}」給 {agent.name}")
try:
result = await agent.execute(task, self.shared_context)
task.result = result
task.status = TaskStatus.DONE
task.finished_at = __import__("time").time()
# 把結(jié)果存入共享上下文
self.shared_context[f"{task.assignee}_result"] = result
self.log(f"? 任務(wù)「{task.name}」完成({task.elapsed_ms:.0f}ms)")
except Exception as e:
task.status = TaskStatus.FAILED
task.error = str(e)
self.log(f"? 任務(wù)「{task.name}」失?。簕e}")
async def run(self, user_request: str) -> dict:
"""主流程:規(guī)劃 → 調(diào)度 → 執(zhí)行 → 匯總"""
print(f"\n{'='*60}")
print(f"?? 用戶(hù)需求:{user_request}")
print(f"{'='*60}\n")
# 1. 任務(wù)規(guī)劃
self.log("開(kāi)始任務(wù)規(guī)劃...")
tasks = await self.plan(user_request)
self.shared_context["user_request"] = user_request
self.task_queue = tasks
print(f"?? 任務(wù)計(jì)劃({len(tasks)} 個(gè)子任務(wù)):")
for i, t in enumerate(tasks, 1):
deps = f"(依賴(lài):{', '.join(t.depends_on)})" if t.depends_on else ""
print(f" {i}. [{t.assignee}] {t.name} {deps}")
print()
# 2. 拓?fù)渑判?+ 執(zhí)行(支持并行)
completed = set()
while len(completed) < len(tasks):
# 找出當(dāng)前可執(zhí)行的任務(wù)(依賴(lài)已完成)
ready = [
t for t in tasks
if t.status == TaskStatus.PENDING
and all(dep in completed for dep in t.depends_on)
]
if not ready:
# 檢查是否死鎖
pending = [t for t in tasks if t.status == TaskStatus.PENDING]
if pending:
self.log("?? 檢測(cè)到任務(wù)依賴(lài)死鎖,強(qiáng)制執(zhí)行剩余任務(wù)")
ready = pending[:1]
else:
break
# 并行執(zhí)行所有可執(zhí)行任務(wù)
await asyncio.gather(*[self._execute_task(t) for t in ready])
completed.update(t.name for t in ready if t.status == TaskStatus.DONE)
# 3. 匯總結(jié)果
return await self._summarize()
async def _summarize(self) -> dict:
"""匯總所有 Agent 的工作成果"""
self.log("匯總所有成果...")
results = {
name: ctx
for name, ctx in self.shared_context.items()
if name.endswith("_result")
}
summary_prompt = f"""
團(tuán)隊(duì)已完成以下工作,請(qǐng)綜合匯總成一份清晰的最終報(bào)告:
{__import__("json").dumps(
{k: str(v)[:500] for k, v in results.items()},
ensure_ascii=False, indent=2
)}
輸出格式:
## 項(xiàng)目總結(jié)
- 完成情況
- 核心交付物
- 關(guān)鍵發(fā)現(xiàn)
- 建議后續(xù)行動(dòng)
"""
summary = await self.think(summary_prompt)
return {
"summary": summary,
"details": results,
"tasks": [t.to_dict() for t in self.task_queue]
}五、實(shí)戰(zhàn):自動(dòng)生成競(jìng)品分析報(bào)告
# main.py
import asyncio
import os
from agents.boss import BossAgent
from agents.researcher import ResearcherAgent
from agents.engineer import EngineerAgent
from agents.tester import TesterAgent
from agents.documenter import DocumenterAgent
API_KEY = os.getenv("OPENAI_API_KEY")
BASE_URL = os.getenv("BASE_URL", "https://api.openai.com/v1")
def build_team() -> dict:
"""組建 AI 工作團(tuán)隊(duì)"""
shared_kwargs = dict(
api_key=API_KEY,
base_url=BASE_URL,
model="gpt-4o",
verbose=True
)
return {
"researcher": ResearcherAgent(**shared_kwargs),
"engineer": EngineerAgent(**shared_kwargs),
"tester": TesterAgent(**shared_kwargs),
"documenter": DocumenterAgent(**shared_kwargs),
}
async def run_competitive_analysis():
"""實(shí)戰(zhàn)場(chǎng)景:自動(dòng)生成競(jìng)品分析報(bào)告"""
team = build_team()
boss = BossAgent(team=team, api_key=API_KEY, base_url=BASE_URL, model="gpt-4o")
result = await boss.run(
"分析 RAG(檢索增強(qiáng)生成)領(lǐng)域的主流工具:LangChain、LlamaIndex、Dify,"
"對(duì)比它們的特點(diǎn)、優(yōu)劣勢(shì)和適用場(chǎng)景,給出選型建議,并用 Python 寫(xiě)一個(gè)簡(jiǎn)單的對(duì)比評(píng)測(cè)腳本。"
)
# 輸出結(jié)果
print("\n" + "="*60)
print("?? 最終報(bào)告")
print("="*60)
print(result["summary"])
# 保存到文件
import json
with open("output/competitive_analysis.json", "w", encoding="utf-8") as f:
json.dump(result, f, ensure_ascii=False, indent=2)
print("\n? 完整結(jié)果已保存到 output/competitive_analysis.json")
async def run_code_generation():
"""實(shí)戰(zhàn)場(chǎng)景:自動(dòng)寫(xiě)代碼 + 測(cè)試 + 文檔"""
team = build_team()
boss = BossAgent(team=team, api_key=API_KEY, base_url=BASE_URL, model="gpt-4o")
result = await boss.run(
"用 Python 實(shí)現(xiàn)一個(gè)高性能的內(nèi)存緩存庫(kù),支持:LRU/LFU/TTL 三種淘汰策略、"
"線(xiàn)程安全、最大容量限制、序列化存儲(chǔ),提供裝飾器 API。"
)
print(result["summary"])
if __name__ == "__main__":
asyncio.run(run_competitive_analysis())
六、Agent 通信與狀態(tài)共享
# 實(shí)現(xiàn)一個(gè)簡(jiǎn)單的消息總線(xiàn),讓 Agent 之間可以通信
import asyncio
from collections import defaultdict
class MessageBus:
"""
Agent 消息總線(xiàn)(發(fā)布-訂閱模式)
允許 Agent 在執(zhí)行過(guò)程中互相通知
"""
def __init__(self):
self._subscribers: dict[str, list[asyncio.Queue]] = defaultdict(list)
def subscribe(self, topic: str) -> asyncio.Queue:
"""訂閱某個(gè)話(huà)題"""
q = asyncio.Queue()
self._subscribers[topic].append(q)
return q
async def publish(self, topic: str, message: dict):
"""發(fā)布消息到某個(gè)話(huà)題"""
for q in self._subscribers[topic]:
await q.put(message)
async def request_help(
self,
from_agent: str,
to_agent: str,
question: str
) -> str:
"""Agent 向另一個(gè) Agent 請(qǐng)求幫助"""
await self.publish(f"help_request_{to_agent}", {
"from": from_agent,
"question": question
})
# 等待響應(yīng)(超時(shí) 30 秒)
response_queue = self.subscribe(f"help_response_{from_agent}")
try:
response = await asyncio.wait_for(response_queue.get(), timeout=30)
return response.get("answer", "")
except asyncio.TimeoutError:
return f"? {to_agent} 未在超時(shí)內(nèi)響應(yīng)"
# 在 Agent 中使用消息總線(xiàn)
class CollaborativeAgent(BaseAgent):
def __init__(self, bus: MessageBus, **kwargs):
super().__init__(**kwargs)
self.bus = bus
async def ask_colleague(self, agent_name: str, question: str) -> str:
"""向同事請(qǐng)教問(wèn)題"""
self.log(f"向 {agent_name} 求助:{question[:50]}...")
answer = await self.bus.request_help(self.name, agent_name, question)
return answer
七、錯(cuò)誤恢復(fù)與重試機(jī)制
import functools
def with_retry(max_retries: int = 3, delay: float = 1.0, backoff: float = 2.0):
"""任務(wù)執(zhí)行失敗自動(dòng)重試的裝飾器(指數(shù)退避)"""
def decorator(func):
@functools.wraps(func)
async def wrapper(self, task: Task, *args, **kwargs):
last_error = None
for attempt in range(max_retries):
try:
return await func(self, task, *args, **kwargs)
except Exception as e:
last_error = e
wait = delay * (backoff ** attempt)
self.log(f"?? 第 {attempt+1} 次失?。▄e}),{wait:.1f}s 后重試...")
await asyncio.sleep(wait)
task.retries += 1
task.status = TaskStatus.FAILED
task.error = str(last_error)
raise last_error
return wrapper
return decorator
# 在 Agent 中使用
class RobustEngineerAgent(EngineerAgent):
@with_retry(max_retries=3, delay=2.0)
async def execute(self, task: Task, shared_context: dict) -> dict:
return await super().execute(task, shared_context)
八、性能優(yōu)化:并行執(zhí)行
# 使用信號(hào)量控制并發(fā),避免 API 限流
import asyncio
from contextlib import asynccontextmanager
class ConcurrentBossAgent(BossAgent):
"""支持并行執(zhí)行的 Boss Agent"""
def __init__(self, *args, max_concurrent: int = 3, **kwargs):
super().__init__(*args, **kwargs)
self.semaphore = asyncio.Semaphore(max_concurrent) # 最多 3 個(gè)并發(fā)
async def _execute_task_with_limit(self, task: Task) -> None:
async with self.semaphore: # 拿到信號(hào)量才執(zhí)行
await self._execute_task(task)
async def run(self, user_request: str) -> dict:
"""并行執(zhí)行不依賴(lài)彼此的任務(wù)"""
tasks = await self.plan(user_request)
# 按依賴(lài)關(guān)系分層
layers = self._topological_sort(tasks)
for layer in layers:
# 同一層的任務(wù)沒(méi)有依賴(lài),可以并行
print(f"? 并行執(zhí)行 {len(layer)} 個(gè)任務(wù):{[t.name for t in layer]}")
await asyncio.gather(*[
self._execute_task_with_limit(t) for t in layer
])
return await self._summarize()
def _topological_sort(self, tasks: list[Task]) -> list[list[Task]]:
"""拓?fù)渑判颍祷乜刹⑿械娜蝿?wù)層"""
task_map = {t.name: t for t in tasks}
layers = []
remaining = set(t.name for t in tasks)
while remaining:
# 找出當(dāng)前沒(méi)有未完成依賴(lài)的任務(wù)
layer = [
task_map[name] for name in remaining
if all(dep not in remaining for dep in task_map[name].depends_on)
]
if not layer:
break # 避免死循環(huán)
layers.append(layer)
remaining -= {t.name for t in layer}
return layers
九、Web 監(jiān)控面板(Gradio)
# monitor.py
import gradio as gr
import asyncio
import json
from main import build_team, BossAgent
task_log = []
def run_agents(user_request: str):
"""在 Gradio 中運(yùn)行多智能體系統(tǒng)"""
team = build_team()
boss = BossAgent(team=team, api_key=os.getenv("OPENAI_API_KEY"))
result = asyncio.run(boss.run(user_request))
task_log.extend(result["tasks"])
# 格式化任務(wù)執(zhí)行結(jié)果
task_table = [[
t["id"], t["name"], t["assignee"],
t["status"], f"{t['elapsed_ms']:.0f}ms"
] for t in result["tasks"]]
return (
result["summary"],
task_table,
json.dumps(result["details"], ensure_ascii=False, indent=2)
)
with gr.Blocks(title="?? 多智能體監(jiān)控面板") as demo:
gr.Markdown("# ?? Python 多智能體協(xié)作系統(tǒng)")
gr.Markdown("輸入任務(wù),5 個(gè) AI Agent 自動(dòng)分工完成,實(shí)時(shí)查看進(jìn)度")
with gr.Row():
user_input = gr.Textbox(
label="?? 輸入任務(wù)",
placeholder="例如:分析 LangChain vs LlamaIndex,給出技術(shù)選型建議...",
lines=3
)
run_btn = gr.Button("?? 啟動(dòng)團(tuán)隊(duì)", variant="primary", size="lg")
summary_output = gr.Markdown(label="?? 最終匯總")
with gr.Accordion("?? 任務(wù)執(zhí)行詳情", open=False):
task_table = gr.Dataframe(
headers=["任務(wù)ID", "任務(wù)名", "負(fù)責(zé)Agent", "狀態(tài)", "耗時(shí)"],
label="任務(wù)執(zhí)行記錄"
)
with gr.Accordion("?? 原始輸出", open=False):
raw_output = gr.JSON(label="詳細(xì)結(jié)果")
gr.Examples(
examples=[
["分析 FastAPI vs Django vs Flask,給出 2026 年 Python API 框架選型建議"],
["用 Python 實(shí)現(xiàn)一個(gè)支持并發(fā)的下載器,帶進(jìn)度條和斷點(diǎn)續(xù)傳"],
["調(diào)研向量數(shù)據(jù)庫(kù) Milvus/Weaviate/Chroma 的優(yōu)劣勢(shì),并寫(xiě)一個(gè)性能對(duì)比腳本"],
],
inputs=user_input
)
run_btn.click(
fn=run_agents,
inputs=user_input,
outputs=[summary_output, task_table, raw_output]
)
demo.launch(server_name="0.0.0.0", server_port=7860)
十、效果對(duì)比:?jiǎn)?Agent vs 多 Agent
| 維度 | 單 Agent | 多 Agent 協(xié)作 |
|---|---|---|
| 輸出質(zhì)量 | 中等,容易混亂 | 高,專(zhuān)業(yè)分工 |
| 復(fù)雜任務(wù) | 容易遺漏、錯(cuò)亂 | 拆解清晰,覆蓋全面 |
| 執(zhí)行速度 | 順序,慢 | 并行,快 2-5x |
| 可擴(kuò)展性 | 差,難以增加能力 | 好,加 Agent 即可 |
| 可調(diào)試性 | 難,輸出混在一起 | 易,每個(gè) Agent 獨(dú)立可觀(guān)測(cè) |
| 適用場(chǎng)景 | 簡(jiǎn)單單步任務(wù) | 復(fù)雜多步驟任務(wù) |
總結(jié)
本文搭建的多智能體系統(tǒng)核心要點(diǎn):
- BaseAgent 基類(lèi):統(tǒng)一的 LLM 調(diào)用 + 工具注冊(cè) + 記憶管理
- 5 個(gè)專(zhuān)業(yè) Agent:研究員/工程師/測(cè)試/文檔/Boss 各司其職
- Boss Agent:任務(wù)規(guī)劃 + 拓?fù)湔{(diào)度 + 結(jié)果匯總
- 并行執(zhí)行:信號(hào)量控制,最大化利用 API 并發(fā)
- 錯(cuò)誤恢復(fù):指數(shù)退避重試,提高任務(wù)成功率
- 消息總線(xiàn):Agent 之間實(shí)時(shí)通信
- Gradio 面板:可視化監(jiān)控執(zhí)行過(guò)程
以上就是Python實(shí)現(xiàn)多智能體協(xié)作的實(shí)戰(zhàn)教程的詳細(xì)內(nèi)容,更多關(guān)于Python多智能體的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
Python基于更相減損術(shù)實(shí)現(xiàn)求解最大公約數(shù)的方法
這篇文章主要介紹了Python基于更相減損術(shù)實(shí)現(xiàn)求解最大公約數(shù)的方法,簡(jiǎn)單說(shuō)明了更相減損術(shù)的概念、原理并結(jié)合Python實(shí)例形式分析了基于更相減損術(shù)實(shí)現(xiàn)求解最大公約數(shù)的相關(guān)操作技巧與注意事項(xiàng),需要的朋友可以參考下2018-04-04
python實(shí)現(xiàn)快速排序的示例(二分法思想)
本篇文章主要介紹了python實(shí)現(xiàn)快速排序的示例(二分法思想),小編覺(jué)得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過(guò)來(lái)看看吧2018-03-03
Python使用holiday-cn庫(kù)實(shí)現(xiàn)節(jié)假日智能識(shí)別工具
holiday-cn 是一個(gè)專(zhuān)門(mén)用于處理中國(guó)節(jié)假日識(shí)別的 Python 庫(kù),它通過(guò)精準(zhǔn)的數(shù)據(jù)模型和高效的判斷邏輯,為開(kāi)發(fā)者提供了一套完整的節(jié)假日處理解決方案,本文通過(guò)代碼示例介紹的非常詳細(xì),需要的朋友可以參考下2026-03-03
pytest?fixtures函數(shù)及測(cè)試函數(shù)的參數(shù)化解讀
這篇文章主要介紹了pytest?fixtures函數(shù)及測(cè)試函數(shù)的參數(shù)化解讀,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-05-05
解決django接口無(wú)法通過(guò)ip進(jìn)行訪(fǎng)問(wèn)的問(wèn)題
這篇文章主要介紹了解決django接口無(wú)法通過(guò)ip進(jìn)行訪(fǎng)問(wèn)的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧2020-03-03
python的Crypto模塊實(shí)現(xiàn)AES加密實(shí)例代碼
這篇文章主要介紹了python的Crypto模塊實(shí)現(xiàn)AES加密實(shí)例代碼,簡(jiǎn)單介紹了實(shí)現(xiàn)步驟,小編覺(jué)得還是挺不錯(cuò)的,具有一定借鑒價(jià)值,需要的朋友可以參考下2018-01-01
python如何通過(guò)實(shí)例方法名字調(diào)用方法
這篇文章主要為大家詳細(xì)介紹了python如何通過(guò)實(shí)例方法名字調(diào)用方法,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2018-03-03

