Token导航 LogoToken导航TokenDH.com

手把手构建企业级 Agent 框架(五):多智能体并行与任务委派,突破单 Agent 的性能与上下文瓶颈

更新时间 2026-05-25来源 Aike正文 1万字阅读约 33分钟2 张图片

在前四篇文章中,我们构建了强大的单 Agent:它拥有 Gateway 多渠道接入、ReAct 运行时动态 Skill 系统。然而,当面对“同时处理数据分析 + 发送报告”或“为不同用户并行执行私有任务”等企业场景时,单一 Agent 的上下文窗口和顺序执行模型就会成为瓶颈。多智能体协作正是破局之道。OpenClaw 通过 Subagent(子代理) 机制优雅地解决了这一问题,而我们将在企业框架中实现一个完整的 Commander-Worker 多智能体系统

系列文章:

手把手构建企业级 Agent 框架:从 OpenClaw 架构到自主实现

手把手构建企业级 Agent 框架(二):Gateway 网关与多渠道接入

手把手构建企业级 Agent 框架(三):Pi Agent 运行时与 ReAct 循环

手把手构建企业级 Agent 框架(四):Skill 系统——知识注入与能力扩展

一、为什么需要 Subagent?单 Agent 的天然局限

即使再强大的 LLM,也面临两个硬性约束:

  • 上下文窗口有限:
    当任务需要查阅大量文档、调用多个工具并处理冗长结果时,单 Agent 的上下文很快就会塞满,导致遗忘或幻觉。
  • 无法并行:
    传统 ReAct 循环是串行的:思考→行动→观察→再思考。如果用户说“同时帮我查销售周报和技术周报”,单 Agent 必须逐一处理,效率极低。

Subagent 的思路非常直觉:将一个大任务拆解为多个子任务,委派给专门的子 Agent 并行执行,然后汇总结果。这就像团队协作:主管(Commander)不需要亲手做所有事,只需分配工作给分析员、工程师,最后整合汇报。

企业级 Subagent 的核心优势:

  • 上下文隔离:
    每个子 Agent 拥有独立的上下文窗口,不会互相污染。
  • 并行加速:
    多个子任务同时执行,总延迟接近最慢的那个子任务。
  • 专业化:
    可以为不同领域定制专用 Agent(如“审批 Agent”“报表 Agent”),配备不同的 Skill 和 Tool。
  • 安全边界:
    敏感操作可以委派给权限更高的子 Agent,实现最小权限原则。

二、Commander-Worker 模式与架构设计

我们采用经典的 Commander(指挥官) + Worker(工人) 拓扑。Commander 就是现有的主 Agent,它通过一个特殊的工具 delegate_to_subagent 来发起委派。Worker 则是独立运行的 Agent 实例,监听任务队列。

图片

工作流程:

  1. 用户发送请求:“请同时生成销售周报和技术周报,并发送给相关人员”。
  2. Commander 的 ReAct 循环在 Planning 阶段分析任务,决定调用两次 delegate_to_subagent,分别委派给 report-sales 和 report-tech Worker,并指定任务描述。
  3. 每次委派调用时,如果允许并行,Commander 不等待结果立即继续(或异步等待)。我们实现一个异步委派,当 Commander 需要结果时再进行汇聚。
  4. 各 Worker 独立执行,返回结果后暂存在注册中心。Commander 最后调用一次 collect_results 工具获取所有结果,并整合回复。

三、接口与协议设计

3.1 子 Agent 注册协议

每个 Worker 需向注册中心声明自己的能力描述、所需工具列表和 Skill 集合。这样 Commander 的 LLM 才能知道何时委派给谁。

from dataclasses import dataclass, field
from typing import List, Dict

@dataclass
classSubagentInfo:
    name:str
    description:str# 一句话说明能力,供 LLM 匹配
    capabilities: List[str]# 能力标签,如 ['report', 'data_analysis']
    required_tools: List[str]
    skill_names: List[str]
    endpoint:str=""# 预留网络调用地址

3.2 任务委派消息

@dataclass
classSubagentTask:
    task_id:str
    agent_name:str
    description:str# 自然语言任务描述
    context: Dict             # 可选的附加上下文(用户ID、租户等)

3.3 SubagentPool 核心接口

classSubagentPool:
defregister(self, info: SubagentInfo, worker:'AgentRuntime'):...
asyncdefdelegate(self, task: SubagentTask)->str:...
asyncdefcollect(self, task_ids: List[str])-> Dict[str,str]:...
defget_all_info(self)-> List[SubagentInfo]:...

四、代码实战:构建可运行的多 Agent 协作系统

我们将扩展上一篇的 agent.py,在主 Agent 的工具列表中注册 delegate_to_subagent 和 collect_subagent_results 两个特殊工具,并启动多个 Worker 实例。

4.1 项目结构

eclaw-subagent/
├── agent.py           # AgentRuntime(含 Commander 逻辑)
├── tools.py           # ToolRegistry + 委派工具实现
├── subagent_pool.py   # SubagentPool 与注册中心
├── worker.py          # 启动 Worker 的脚本
├── main.py            # 测试启动入口
├── memory.py          # 记忆(沿用)
├── skill.py           # Skill 加载器(沿用)
└── skills/            # 技能文件

4.2 SubagentPool 实现 (subagent_pool.py)

import asyncio
import uuid
from typing import Dict, List
from dataclasses import dataclass, field

@dataclass
classSubagentInfo:
    name:str
    description:str
    capabilities: List[str]= field(default_factory=list)
    required_tools: List[str]= field(default_factory=list)
    skill_names: List[str]= field(default_factory=list)

@dataclass
classSubagentTask:
    task_id:str
    agent_name:str
    description:str
    context: Dict = field(default_factory=dict)

classSubagentPool:
def__init__(self):
        self._workers: Dict[str,'AgentRuntime']={}# agent_name -> worker实例
        self._info: Dict[str, SubagentInfo]={}
        self._results: Dict[str,str]={}# task_id -> result
        self._lock = asyncio.Lock()

defregister(self, info: SubagentInfo, worker:'AgentRuntime'):
        self._workers[info.name]= worker
        self._info[info.name]= info
print(f"[SubagentPool] 注册 Worker: {info.name} - {info.description}")

defget_all_info(self)-> List[SubagentInfo]:
returnlist(self._info.values())

asyncdefdelegate(self, task: SubagentTask)->str:
"""
        异步委派任务,立即返回 task_id。
        实际任务在后台执行,结果存入 self._results。
        """
        worker = self._workers.get(task.agent_name)
ifnot worker:
returnf"错误:Agent '{task.agent_name}' 未注册"
# 启动后台任务
        asyncio.create_task(self._execute_worker(worker, task))
returnf"任务已委派,task_id={task.task_id}"

asyncdef_execute_worker(self, worker:'AgentRuntime', task: SubagentTask):
"""执行 Worker 并存储结果"""
# 为 Worker 构造一个简单的输入消息
        result_chunks =[]
asyncfor event in worker.run(
            session_id=f"subagent-{task.task_id}",
            user_input=task.description,
            ctx=None# 可传入安全上下文
):
if event["type"]=="final":
                result_chunks.append(event["data"])
        final_result ="\n".join(result_chunks)if result_chunks else"(无结果)"
asyncwith self._lock:
            self._results[task.task_id]= final_result

asyncdefcollect(self, task_ids: List[str])-> Dict[str,str]:
"""等待并收集指定 task_id 的结果(可设置超时)"""
        results ={}
for tid in task_ids:
# 简单轮询等待(生产环境可用 asyncio.Event)
for _ inrange(30):# 最多等 30 秒
if tid in self._results:
                    results[tid]= self._results.pop(tid)
break
await asyncio.sleep(0.3)
else:
                results[tid]="超时未返回结果"
return results

4.3 将委派工具注册到 ToolRegistry

在 tools.py 中新增两个工具函数,并将它们的引用传入 AgentRuntime。由于工具需要访问 SubagentPool,我们通过闭包或部分应用方式绑定。

# tools.py 中添加
from subagent_pool import SubagentPool, SubagentTask
import uuid

defregister_delegation_tools(registry, pool: SubagentPool):
asyncdefdelegate_to_subagent(params:dict, ctx)->dict:
        agent_name = params["agent_name"]
        description = params["task_description"]
        task_id =str(uuid.uuid4())
        task = SubagentTask(task_id=task_id, agent_name=agent_name, description=description)
        result =await pool.delegate(task)
return{"task_id": task_id,"status": result}

asyncdefcollect_subagent_results(params:dict, ctx)->dict:
        task_ids = params["task_ids"]
        results =await pool.collect(task_ids)
return results

from dataclasses import dataclass as td
    registry.register(ToolDefinition(
        name="delegate_to_subagent",
        description="将子任务委派给指定的子 Agent。参数:agent_name (子Agent名称), task_description (任务描述)。返回 task_id。",
        parameters_schema={
"type":"object",
"properties":{
"agent_name":{"type":"string"},
"task_description":{"type":"string"}
},
"required":["agent_name","task_description"]
},
        fn=delegate_to_subagent
))
    registry.register(ToolDefinition(
        name="collect_subagent_results",
        description="收集已委派子任务的结果。参数:task_ids (task_id 列表)。",
        parameters_schema={
"type":"object",
"properties":{
"task_ids":{"type":"array","items":{"type":"string"}}
},
"required":["task_ids"]
},
        fn=collect_subagent_results
))

4.4 修改 AgentRuntime 初始化以支持 SubagentPool

# agent.py 中的 AgentRuntime 修改 __init__
classAgentRuntime:
def__init__(self,... subagent_pool: SubagentPool =None):
        self.tools = ToolRegistry()
        self.tools.register_default_tools()
if subagent_pool:
            register_delegation_tools(self.tools, subagent_pool)
# ...其余不变

4.5 构建专用 Worker 实例

我们可以为不同 Worker 定制 Skill 集合和工具。例如报表 Agent 拥有生成报告的 Skill,审批 Agent 拥有审批相关工具。

在 worker.py 中创建 Worker 实例并注册到 pool:

# worker.py
from subagent_pool import SubagentPool, SubagentInfo
from agent import AgentRuntime  # 复用同一个类
import asyncio

asyncdefsetup_workers(pool: SubagentPool):
# 创建一个报表 Worker
    report_agent = AgentRuntime()
# 可以动态注入特定的技能(假设 SkillLoader 加载了 report_skill.md)
    report_agent.skills._load_all()
    pool.register(SubagentInfo(
        name="report-worker",
        description="擅长生成各类业务报表,如销售周报、技术周报",
        capabilities=["report","data_summary"],
        required_tools=["query_order"],# 简化
        skill_names=["report_skill"]
), report_agent)

# 创建一个审批 Worker
    approval_agent = AgentRuntime()
    pool.register(SubagentInfo(
        name="approval-worker",
        description="处理退款审批、权限申请等流程",
        capabilities=["approval","workflow"],
        required_tools=["send_notification"],
        skill_names=["approval_skill"]
), approval_agent)

4.6 主测试入口 (main.py)

我们将 Commander 与 Worker 放在同一个进程中运行,通过 asyncio 协作。

# main.py
import asyncio
from subagent_pool import SubagentPool
from worker import setup_workers
from agent import AgentRuntime, AgentContext

asyncdefmain():
    pool = SubagentPool()
await setup_workers(pool)

# 创建 Commander(主 Agent),并注入 pool
    commander = AgentRuntime(subagent_pool=pool)
# 为 Commander 也配置一些技能(例如如何委派)
    ctx = AgentContext(request=None, security_roles=["admin"])

# 模拟用户请求
    user_input ="请同时帮我生成销售周报和技术周报,完成后告诉我结果。"

print("=== 用户请求 ===")
print(user_input)
print("\n=== Agent 执行流 ===")
asyncfor event in commander.run("session-main", user_input, ctx):
if event["type"]=="thought":
print(f"[思考] {event['data']}")
elif event["type"]=="tool_call":
print(f"[调用工具] {event['data']}")
elif event["type"]=="tool_result":
print(f"[工具结果] {event['data']}")
elif event["type"]=="final":
print(f"[最终回复] {event['data']}")

if __name__ =="__main__":
    asyncio.run(main())

4.7 模拟 LLM 决策逻辑(改动 agent.py 中的 _mock_llm)

为了让 Commander 知晓何时使用委派工具,我们需要更新模拟 LLM,使其在检测到“同时生成…周报”之类的并行任务时,返回 多次工具调用 或 依次委派。这里简化为在 _mock_llm 中检测关键词“同时”并返回委派指令。真实场景中 LLM 会自行规划。

# agent.py 中 _mock_llm 的增强
asyncdef_mock_llm(self, messages: List[Dict])-> Dict:
    last_user = messages[-1]["content"]if messages else""
# 检查是否已有委派指示(避免循环)
# 简单示例:如果用户请求中包含“同时”,则依次委派。
if"同时"in last_user and"周报"in last_user:
# 第一次返回委派销售周报,如果历史中已有委派则进行第二次
        has_delegated_sales =any("delegate_to_subagent"instr(m)for m in messages)
ifnot has_delegated_sales:
return{"decision":{
"thought":"用户需要并行生成两份周报,我先委派销售周报任务。",
"action":"tool_call",
"tool_name":"delegate_to_subagent",
"tool_params":{"agent_name":"report-worker","task_description":"生成销售周报"}
}}
else:
# 检查是否已委派技术周报
            has_delegated_tech ="技术周报"instr(messages)
ifnot has_delegated_tech:
return{"decision":{
"thought":"接着委派技术周报任务。",
"action":"tool_call",
"tool_name":"delegate_to_subagent",
"tool_params":{"agent_name":"report-worker","task_description":"生成技术周报"}
}}
else:
# 两次委派后,收集结果
# 需要解析之前的 task_id(简化,硬编码)
return{"decision":{
"thought":"两项任务均已委派,现在收集结果。",
"action":"tool_call",
"tool_name":"collect_subagent_results",
"tool_params":{"task_ids":["task-sales","task-tech"]}
}}
# ...原有逻辑

⚠️ 注意: 上述模拟 LLM 仅为演示多步委派流程。实际接入 GPT-4 等强力模型时,LLM 会自动根据工具描述和用户请求生成正确的调用序列,无需硬编码判断。

五、运行与测试

启动 main.py,观察 Commander 如何委派并收集结果。预期输出类似:

=== Agent 执行流 ===
[思考] 用户需要并行生成两份周报,我先委派销售周报任务。
[调用工具] {'tool': 'delegate_to_subagent', 'params': {'agent_name': 'report-worker', 'task_description': '生成销售周报'}}
[工具结果] {'task_id': 'abc123...', 'status': '任务已委派...'}
[思考] 接着委派技术周报任务。
[调用工具] {'tool': 'delegate_to_subagent', 'params': {'agent_name': 'report-worker', 'task_description': '生成技术周报'}}
[工具结果] {'task_id': 'def456...', 'status': '任务已委派...'}
[思考] 两项任务均已委派,现在收集结果。
[调用工具] {'tool': 'collect_subagent_results', 'params': {'task_ids': ['abc123', 'def456']}}
[工具结果] {'abc123': '销售周报已生成:...', 'def456': '技术周报已生成:...'}
[最终回复] 销售周报和技术周报已完成,内容分别是...

Worker 实例在后台实际执行了各自的 Agent 循环,由于我们的 Worker 也使用 AgentRuntime,它们可以调用自己注册的工具(如 query_order 或 send_notification)并注入各自的 Skill 指导。

六、生产化考虑:消息队列与网络分离

上述示例将 Commander 和 Worker 放在同进程只是便于演示。在生产环境中,你应该将 SubagentPool 替换为基于消息队列(Redis/Kafka)的实现,使得 Worker 可以运行在独立的服务器上,支持弹性伸缩。接口保持 delegate 和 collect 不变即可。

📌 扩展建议:

  • 为每个 Worker 添加健康检查,注册中心定期心跳检测。
  • 实现任务优先级和队列长度限制。
  • 记录完整的调用链:Commander → 哪个 Worker,耗时多少,便于监控。

七、与 OpenClaw 的对标思考

🔍 “我们做了什么” vs “OpenClaw 为什么这样做”?

拓扑结构:OpenClaw 的 Subagent 也采用主从模式,主 Agent 通过announce机制等待结果。我们的delegate+collect异曲同工,但显式分离了委派和收集阶段,更易于实现并行。

通信机制:OpenClaw 内部使用 Node.js 的 EventEmitter 或文件系统进行进程间通信。我们选择用 asyncio 任务在单进程内模拟,并预留了网络接口,适合向微服务架构演进。

能力注册:OpenClaw 的 Worker 自动暴露自身能力(通过 Skill 列表),我们的实现手动注册SubagentInfo,这种显式声明在企业管理中更可控,也便于加入权限验证。

错误处理:我们的collect中加入超时机制,避免了无限等待。企业生产中还应加入重试和死信处理。

总的来说,我们的 Subagent 系统继承了 OpenClaw 主子协作的先进思想,同时在可观测性、安全和跨网络部署方面进行了企业级增强。

八、总结与下一步

本文我们:

  1. 深入分析了单 Agent 瓶颈,并引入 Commander-Worker 模式来解决并行与上下文隔离。
  2. 设计了 SubagentPool 注册中心,实现了任务委派、异步执行和结果收集。
  3. 将委派与收集包装为两个标准工具,使主 Agent 的 ReAct 循环能自主决定何时使用。
  4. 编写了完整的测试入口,展示了并行委派流程。

下一篇文章预告:《第6篇:双源记忆系统》——我们将为 Agent 装上长短期记忆,通过 短期上下文窗口 + 长期向量记忆 实现跨会话的知识保留,距离真正的智能同事又进一步。

图片
文章标签智能体
资讯来源:由AI资讯编辑整理自互联网公开内容,版权归原作者所有,未经许可,不得转载。

继续浏览更多资讯

返回资讯目录

相关资讯

更多