最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Python實(shí)現(xiàn)多智能體協(xié)作的實(shí)戰(zhàn)教程

 更新時(shí)間:2026年04月24日 09:08:49   作者:lbb?小魔仙  
一個(gè) Agent 解決不了的問(wèn)題,就組一個(gè)隊(duì),本文將帶大家使用 Python 從零搭建多智能體協(xié)作系統(tǒng),文中的示例代碼講解詳細(xì),感興趣的小伙伴可以了解下

一個(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]

要求:

  1. 覆蓋所有公開(kāi)函數(shù)/方法
  2. 包含:正常用例、邊界條件、異常用例
  3. 使用 pytest.mark.parametrize 參數(shù)化測(cè)試
  4. 添加 Mock(如果涉及外部依賴(lài))
  5. 目標(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ù)的方法

    這篇文章主要介紹了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)SOM算法

    python實(shí)現(xiàn)SOM算法

    這篇文章主要為大家詳細(xì)介紹了python實(shí)現(xiàn)SOM算法,聚類(lèi)算法,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-02-02
  • python實(shí)現(xiàn)快速排序的示例(二分法思想)

    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í)別工具

    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ù)化解讀

    這篇文章主要介紹了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)題

    這篇文章主要介紹了解決django接口無(wú)法通過(guò)ip進(jìn)行訪(fǎng)問(wèn)的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-03-03
  • 在PyCharm中實(shí)現(xiàn)添加快捷模塊

    在PyCharm中實(shí)現(xiàn)添加快捷模塊

    今天小編就為大家分享一篇在PyCharm中實(shí)現(xiàn)添加快捷模塊,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。一起跟隨小編過(guò)來(lái)看看吧
    2020-02-02
  • 如何用Python從桌面讀取二維碼信息詳解

    如何用Python從桌面讀取二維碼信息詳解

    二維碼作為一種信息傳遞的工具,在當(dāng)今社會(huì)發(fā)揮了重要作用,下面這篇文章主要給大家介紹了關(guān)于如何用Python從桌面讀取二維碼信息的相關(guān)資料,文中通過(guò)示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2021-10-10
  • python的Crypto模塊實(shí)現(xiàn)AES加密實(shí)例代碼

    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)用方法

    python如何通過(guò)實(shí)例方法名字調(diào)用方法

    這篇文章主要為大家詳細(xì)介紹了python如何通過(guò)實(shí)例方法名字調(diào)用方法,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-03-03

最新評(píng)論

台东市| 富川| 五寨县| 塘沽区| 玉田县| 屏东市| 宜良县| 米易县| 南昌县| 曲沃县| 称多县| 明光市| 康平县| 长武县| 许昌市| 临汾市| 原阳县| 博乐市| 泽库县| 蛟河市| 若羌县| 营口市| 老河口市| 桦川县| 客服| 普格县| 抚顺市| 桐乡市| 保山市| 沧源| 磴口县| 开封市| 简阳市| 高要市| 左云县| 兴仁县| 旬阳县| 庄浪县| 大厂| 贞丰县| 乐安县|