数据截至 (上游 commit 92c146faa529)
多智能体编排:Task、AgentTeam 与三种 Process
30 秒导读: 上一章讲的是"一个 Agent 怎么自己转一圈"。这一章讲"一支团队怎么协作":你把工作拆成若干
Task(工作单元),交给AgentTeam(编排器),再由Process(调度策略)决定这些任务按什么顺序、由哪个 agent、跑几遍。三种策略——顺序sequential、由经理 LLM 临场指挥的hierarchical、带图和条件循环的workflow。
本章的边界:这里讲的是 LLM 自主编排——顺序由代码框架 + LLM 决策共同决定。开发者"写死一张确定性图、每一步都不问 LLM"的那种流水线,是 05 章 AgentFlow 的内容。同组其它章见 index.md。
1. 这是什么(零基础也能懂)
一句话定义: 把"一个 agent"升级成"一支 agent 团队"——你描述好每件事(Task),团队按某种策略把它们跑完。
解决什么问题: 真实任务往往一个 agent 干不完,或者干不好。
- "先研究、再写稿"——这是两步,后一步要吃前一步的结果。
- "让一个懂全局的经理决定现在该派谁"——这是动态调度,顺序事先不知道。
- "对列表里每一项都跑一遍,分数不够就重来"——这是图 + 循环 + 条件。
一个 agent 撑不起这些;需要一个能"排班"的东西。PraisonAI 里那个东西叫 AgentTeam。
用起来什么样: 最小的顺序团队——研究员先跑,作者吃到研究结果再跑。
# 示意,非源码。演示三件套:Agent + Task + AgentTeam
from praisonaiagents import Agent, Task, AgentTeam
researcher = Agent(role="研究员", instructions="研究给定话题")
writer = Agent(role="作者", instructions="根据研究结果写文章")
t1 = Task(description="研究 2026 年 AI agent 趋势", agent=researcher)
t2 = Task(description="写一篇 800 字文章", agent=writer, depends_on=[t1]) # 依赖 t1 的产出
team = AgentTeam(agents=[researcher, writer], tasks=[t1, t2], process="sequential")
result = team.start() # t1 先跑,t2 吃到 t1 的结果再跑
一句话直觉: 把 Task 当"工单",AgentTeam 当"项目组长",Process 当"排班表"。你只写工单和组员;怎么派活,交给排班表。
2. 顶层全景(它大概怎么转)
先看三个部件怎么分工。这张图从上到下是控制流:你调 start(),它委托给 Process 出一个"下一个该跑谁"的序列,AgentTeam 逐个执行。
你的代码
│ team.start() / team.run()
▼
┌───────────────────────────┐
│ AgentTeam(编排器) │ agents/agents.py:554
│ 持有 agents + tasks + 共享state │
└───────────┬───────────────┘
│ run_all_tasks() 按 process 分派
▼
┌───────────────────────────┐
│ Process(调度策略) │ process/process.py:20
│ 生成"下一个 task_id"的序列 │ (generator,yield 一个跑一个)
└───────────┬───────────────┘
yield task_id │ 逐个回传
▼
┌───────────────────────────┐
│ run_task → execute_task │ 带重试 + guardrail 校验
│ 真正调某个 Agent 跑这条 Task│
└───────────────────────────┘
三个部件一句话职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
Task | 一条工作单元:描述、归属 agent、依赖、条件路由、guardrail、记忆 | task/task.py:20 |
AgentTeam | 编排器:持有 agents/tasks/共享状态,驱动执行、管重试与 guardrail | agents/agents.py:620 |
Process | 调度策略:三种算法各自决定"下一个跑哪条 task" | process/process.py:20 |
Handoff | agent 之间的横向交接:把当前对话转交给另一个专职 agent | agent/handoff.py:297 |
主线走一遍(高层):
- 你
AgentTeam(tasks=[...], process="sequential").start()。 start()调run_all_tasks()(agents/agents.py:1889),按self.process分派到Process的对 应方法。Process方法是个 generator:每yield一个task_id,AgentTeam就run_task()跑一条,跑完再回来要下一个。run_task里有重试循环 + guardrail 校验;跑完写状态、存记忆、触发回调。- 全部跑完,
start()默认返回最后一条 task 的产出(return_dict=True才返回全量字典)。
关键设计:编排器和调度策略解耦。AgentTeam 只会"跑一条 task",不关心顺序;顺序全由 Process 的 generator 吐。想换调度策略,只换 process= 字符串即可。
3. 核心机制之一:Task —— 工作单元
它要解决的小问题
把"一件要做的事"变成一个可被调度、可依赖、可路由、可校验、可记忆的对象。不是简单一个字符串描述。
Task 长什么样
Task.__init__(task/task.py:47)的参数多到吓人——因为它是"统一工作单元",既当 AgentTeam 的任务, 又当 Workflow 的步骤,还吸收了历史上多个类的字段(源码注释直说 "unified abstraction … supports all features from the legacy Task class",task/task.py:21-25)。但核心就几组:
| 关切 | 关键字段 | 说明 |
|---|---|---|
| 干什么 | description / action / handler | 三选一必须有其一,否则 __init__ 直接 raise(task/task.py:154)。action 是 description 的别名 |
| 归属 | agent | 由哪个 agent 执行 |
| 依赖 | depends_on / context | 前置任务;depends_on 是 context 的别名,语义更清楚(task/task.py:55) |
| 路由 | condition / when / then_task / else_task | 两套条件路由,见 §6 |
| 校验 | guardrail / guardrails / max_retries | 产出校验 + 重试上限,见 §7 |
| 记忆 | memory / config['memory_config'] | 产出落库,见本节末 |
| 类型 | task_type | "task" / "decision" / "loop",决定 workflow 里怎么路由 |
依赖 = context = 拿到上游产出
depends_on 不只是"排个先后",它真正的作用是把上游 task 的 产出喂进本 task 的 prompt。
depends_on 只是 context 的别名(property 见 task/task.py:424-432)。执行时,框架遍历 context,把每个前置 Task 的 result.raw 拼进当前任务的提示词(execute_callback 里的 prompt 组装,task/task.py:827-851):
# 示意,非源码。展示 context 如何变成 prompt 里的一段
for item in task.context:
if isinstance(item, str):
blocks.append(f"Input Content:\n{item}")
elif hasattr(item, "result") and item.result: # 前置 Task 对象
blocks.append(item.result.raw) # 直接塞进上游产出正文
# 去重后拼成 "Context:\n..." 追加到本任务 prompt
重点看: 依赖关系不是抽象的图边,它落地成"下游能读到上游写了什么"。这就是任务链(task chaining)的物理实现。
状态机:合法流转 + 向后兼容
Task 有 5 个状态(task/protocols.py:20-24):not started / in progress / completed / failed / cancelled。
合法流转写在 _VALID_TRANSITIONS(task/protocols.py:43-47):
not started ──▶ in progress ──▶ completed(终态)
│ │
├──▶ failed ◀────┘
│ │
│ └──▶ in progress(失败可重试)
└──▶ cancelled
set_status(task/task.py:502)有个务实的取舍:非法流转只警告、不拦截。它先问 _lifecycle_manager.can_transition(),不合法就打 warning,然后照样改(task/task.py:514-523,注释 "Allowing for backward compatibility")。设计者选择了"向后兼容 > 状态机严格性"。
任务级记忆:产出可落库
如果给了 memory 或 config['memory_config'],任务产出会存进长期记忆。
initialize_memory(task/task.py:562):双重检查锁懒加载Memory实例(task/task.py:567-572,先无锁判空、拿锁后再判一次),避免并发重复初始化。store_in_memory(task/task.py:657):把内容store_long_term,附上agent_name/task_id/timestamp元数据。失败只记进non_fatal_errors,除非fail_on_memory_error=True才抛(task/task.py:671-677)。
记忆的完整体系是 06 章 的内容,这里只需知道"任务产出能自动进记忆"。
4. 核心机制之二:AgentTeam —— 编排器
它要解决的小问题
拿着一堆 agent 和 task,驱动它们跑完,同时处理三件杂事:重试、guardrail 校验、跨任务共享状态。
两个入口:start vs run
| 方法 | 用途 | 区别 |
|---|---|---|
start() | 交互/调试 | TTY 里默认打 Rich 面板显示进度(agents/agents.py:1964) |
run() | 生产/脚本 | 静默,内部就是 start(output="silent")(agents/agents.py:2192-2204) |
astart() | 异步版 | 走 arun_all_tasks,支持并行跑 async_execution=True 的任务(agents/agents.py:1889) |
两者最终都调 run_all_tasks()。返回值默认是最后一条 task 的 result.raw;传 return_dict=True 才拿到 {task_status, task_results} 全量字典(agents/agents.py:2177-2190)。
分派:run_all_tasks 是个 3 路开关
run_all_tasks(agents/agents.py:1889)本身很薄——就是按 self.process 把活转给 Process 的对应 generator,然后逐个 run_task:
# 真实结构(agents/agents.py:1725-1741),已简化
if self.process == "workflow":
for task_id in process.workflow(): # deprecated,见 §5 说明
self.run_task(task_id)
elif self.process == "sequential":
for task_id in process.sequential():
self.run_task(task_id)
elif self.process == "hierarchical":
for task_id in process.hierarchical(): # 经理可能临场造 task
if isinstance(task_id, Task):
task_id = self.add_task(task_id)
self.run_task(task_id)
注意 hierarchical 分支的特殊处理:经理会 yield 一个新造的 manager Task(不是 id),所以这里要先 add_task 把它登记进去。
执行一条:run_task 的重试循环
run_task(agents/agents.py:1818)是真正的干活循环。骨架:
while 未完成 and 重试次数 < max_retries:
task_output = execute_task(task_id) # 调 agent 真跑一次
if 产出通过 completion_checker:
task_output, should_retry = _apply_task_guardrail(...) # guardrail 校验
if should_retry: 重试++;continue
标记 completed → execute_callback(存记忆) → 存文件 → on_task_complete 回调
else:
标记 in progress → 指数退避 sleep → 重试++
用尽重试仍未完成 → 标记 failed
execute_task(agents/agents.py:1780)本身只做三步(用了 DRY 辅助函数):构造执行上下文 _build_execution_context → _execute_with_agent_sync 真调 agent → _process_task_result 收结果。
guardrail 应用:_apply_task_guardrail
_apply_task_guardrail(agents/agents.py:1280)是"校验 + 重试信号"的关口。它调 task 自己的 _process_guardrail(见 §7),然后:
- 校验失败且还有重试额度 →
retry_count++、状态回in progress、返回(output, True)让run_task重跑(agents/agents.py:1294-1304)。 - 校验失败且重试用尽 → 直接
raise(agents/agents.py:1295-1299)。 - 校验通过且 guardrail 改写了产出 → 用新值替换
task.result(agents/agents.py:1307-1320)。
重点看: 重试的"要不要再来一次"由 guardrail 决定,但"再来一次"这个动作发生在 run_task 的循环里。职责分离:_apply_task_guardrail 只判断并发信号,run_task 执行重跑。
共享状态:team 级黑板
set_state / get_state / update_state / has_state / clear_state(agents/agents.py:2386 起)是一组加锁的键值存储,给整个团队当"黑板":任何 task 都能读写,用于跨任务传递框架级信号。写操作都在 with self._state_lock 里(agents/agents.py:2388),读操作直接读 dict。还有 increment_state(agents/agents.py:2422)做原子自增。
5. 核心机制之三:三种 Process(调度策略)
三种策略都是 Process 类(process/process.py:20)上的方法,每个都是 generator——yield 一个 task_id、等 AgentTeam 跑完、再决定下一个。每种都有 sync 和 async 两个版本(asequential/ahierarchical/aworkflow)。
先一眼看清三者的定位差异:
| 策略 | 顺序由谁定 | 适合 | 复杂度 |
|---|---|---|---|
sequential | 固定:tasks 声明顺序 | 线性流水,A→B→C | 最简单 |
hierarchical | 经理 LLM 每轮临场决定 | 顺序事先不知道、要动态派活 | 中 |
workflow | 图 + 条件 + 循环 | 分支、回环、按结果路由 | 最复杂 |
注意
workflow的现状: 同步的Process.workflow()(process/process.py:1190)已被标记 deprecated,docstring 直接建议改用新的Workflow类(process/process.py:1193-1210)。但它仍是理解"图式编排"的最好样本,且 async 的aworkflow仍在用。本节讲它的原理;确定性的新写法见 05 章。
5.1 sequential —— 按声明顺序跑
最朴素:遍历 self.tasks,没完成的就 yield。
# 真实实现,已简化(process/process.py:1508-1524)
def sequential(self):
for task_id in self.tasks:
task = self.tasks[task_id]
if task.status == "failed":
continue # 跳过已失败的
if 任一 context 依赖已 failed:
task.status = "failed"; continue # 依赖挂了,自己也挂
if task.status != "completed":
yield task_id
一个不显然的细节(process/process.py:1625-1632):依赖失败会传染。如果某 task 的 context 里有 task 已经 failed,它自己直接标 failed、不跑。避免"拿着半截上游数据硬跑"。
5.2 hierarchical —— 经理 LLM 临场指挥
思路: 不预设顺序,而是造一个"经理 agent",每一轮把所有任务的现状喂给它,让它用 LLM 决定"下一个跑哪条、派给谁、还是停"。
hierarchical(process/process.py:1636)的流程:
造 manager_agent(role=项目经理)+ manager_task
│
▼
while 完成数 < 总数:
把所有 task 的 {id,name,status,agent} 汇总成 tasks_summary
问经理:返回 {task_id, agent_name, action} ← ManagerInstructions 结构
├─ action == "stop" → 结束
├─ task_id 非法 → 回灌错误、重问(最多 3 次)
└─ 合法 → 把该 task 改派给 agent_name → yield 执行
经理的输出被约束成一个 Pydantic 结构(process/manager_schema.py:6):
class ManagerInstructions(BaseModel):
task_id: int
agent_name: str
action: str # "execute" 或 "stop"
容错是这里的精华。 让 LLM 出结构化 JSON 天生不可靠,所以有两层兜底:
- 拿结构化输出的双路 fallback ——
_get_manager_instructions_with_fallback(process/process.py:526):先试output_pydantic结构化输出;失败就退到"塞一段 JSON schema 描述进 prompt +output_json=True"再试(process/process.py:543-568);还失败才抛。解析统一走_parse_manager_instructions(process/process.py:185),json.loads+ Pydantic 校验。 - 经理选了不存在的 task 怎么办 —— 不是崩溃,而是把合法 id 列表回灌给经理让它重选,最多
MAX_INVALID_SELECTIONS=3次(process/process.py:1715-1741)。async 版还给经理的解析失败加了指数退避重试(process/process.py:1093-1121)。
选中后,经理还能临场换执行人:遍历 self.agents 找到 agent_name 匹配的,直接改 task.agent(process/process.py:1747-1753)。这是 hierarchical 相比 sequential 的核心增益——动态派活。
5.3 workflow —— 图 + 条件 + 循环
思路: 任务之间用 next_tasks(下一步)和 condition(按结果分支)连成一张有向图,支持回环和"按列表逐项循环"。
workflow/aworkflow(process/process.py:1190 / 591)开跑前先做两件准备:
- 建反向边: 遍历所有 task 的
next_tasks,给目标 task 的previous_tasks追加自己(process/process.py:1223-1228)。这样每条 task 既知道"下一步是谁",也知道"我上一步是谁"。 - 找起点: 找
is_start=True的;没有就用第一条(process/process.py:1231-1239)。
然后进主循环,每轮做这些事(带 max_iter 和 workflow_timeout 双重刹车,process/process.py:1271-1291):
每轮:
① 全部 completed? → 结束
② 构造本任务 context(_build_task_context)存进 _execution_context
③ 若是 loop 任务 → 展开子任务(_create_loop_subtasks)
否则 → yield 执行
④ 完成的普通任务重置回 "not started"(允许再跑)—— 除非 loop/async/rerun=False
⑤ 决定下一个:
decision/loop 任务 → 按 condition[decision] 路由,遇 exit 就结束
否则 → 走 next_tasks[0]
都没有 → _find_next_not_started_task 兜底
四个值得单独看的设计:
① loop 任务 = 把一个文件/列表铺开成一串子任务。 _create_loop_subtasks(process/process.py:204)读 input_file(CSV 或文本),每行造一个子任务,用 next_tasks 把它们串成链(process/process.py:277-281)。CSV 还会把多列拼成 Question/Answer 对(process/process.py:231-234)。这就是"对列表每一项跑一遍"的实现。
② 完成即重置,允许回环。 主循环末尾会把刚完成的普通任务状态重置回 not started(process/process.py:865-891 的 async 版),否则回环时它会被当成"已完成"跳过。但 loop 任务、并行的 async_execution 任务、rerun=False 的任务不重置——避免无限重跑。
③ context 的两种粒度。 _build_task_context(process/process.py:392)按 retain_full_context 选:True 就把所有 previous_tasks 的产出都拼进去;False(默认)只拼最近一条(process/process.py:418-446)。省 token 还是要全景,交给用户配。它还会把上一轮的 validation_feedback 拼进去(process/process.py:395-406),让重试的任务知道"上次为什么被打回"。
④ 卡死兜底。 如果按图走不出下一个任务,_find_next_not_started_task(process/process.py:455)扫一遍找还没跑、且有出边(condition 或 next_tasks)的任务;每个任务有独立重试计数,超过 max_retries 就标 failed(process/process.py:482-491)。防止图有断点时整个流程静默挂起。
6. 条件路由:两套并存的分支系统
Task 能"按结果决定下一步走哪",但 PraisonAI 里有两套条件路由,历史原因并存。看清楚别混:
| 老系统(condition dict) | 新系统(when/then/else) | |
|---|---|---|
| 字段 | task_type="decision"/"loop" + condition={...} | when + then_task + else_task |
| 判断依据 | LLM 产出里的 decision 字段字符串 | 表达式对 context 求值(如 "{{score}} > 80") |
| 谁在用 | Process.workflow/aworkflow 主循环 | Task 级 evaluate_when/get_next_task 辅助方法 |
| 源码 | process/process.py:902-941 | task/task.py:955 / 957 |
老系统怎么路由: decision/loop 任务跑完,从 result.pydantic.decision 或 result.raw 取一个决策字符串,拿它去查 condition 字典(process/process.py:897-914)。condition 形如 {"done": ["next_task"], "retry": ["self"], "exit": []};命中 exit 或空列表就结束工作流(process/process.py:906-910)。decision 任务的 output_pydantic 是动态生成的——用 condition 的键当 Literal 枚举(task/task.py:392-398),强制 LLM 只能从合法分支里选。
新系统怎么路由: evaluate_when(task/task.py:955)把 when 表达式丢给共享的 evaluate_condition(与 AgentFlow 同一套求值器)求真假;get_next_task(task/task.py:980)据此返回 then_task 或 else_task。它是给更新的路由路径用的,和 AgentFlow 的条件语法统一(源码注释 "same as AgentFlow")。
验证反馈闭环(老系统的巧处): 当 decision 判成"失败类"(invalid/retry/reject 等,列在 VALIDATION_FAILURE_DECISIONS,process/process.py:23),框架不只是路由回去,还把"为什么失败、被拒的产出是什么"打包成 validation_feedback 挂到目标任务上(process/process.py:925-937)。下次这个任务跑时,_build_task_context 会把这段反馈拼进 prompt——重试带着"上次错哪了"的信息,而不是盲目再来一遍。
7. Guardrail:任务级产出校验
它要解决的小问题: agent 产出可能不合格(格式错、内容假、没答到点)。Guardrail 是"产出的守门员",不合格就打回重试。
一条 Task 可以带 guardrail,两种形态(task/task.py:78-79,guardrails 复数是正名、guardrail 单数已废弃):
- 可调用函数 —— 签名必须收一个
TaskOutput、返回(bool, Any)或GuardrailResult。_setup_guardrail(task/task.py:434)在构造时就用inspect校验签名和返回注解,不合规直接raise(task/task.py:447-476)。 - 字符串描述 —— 会包成
LLMGuardrail,用 agent 的 LLM 去判"这产出符合这句要求吗"(task/task.py:479-498);此时没有 agent 会报错。
运行时 _process_guardrail(task/task.py:862)调那个函数,统一转成 GuardrailResult;函数内部抛异常也会被兜成一个 success=False 的结果(task/task.py:887-894),不会让校验本身把流程搞崩。
三者怎么串起来,一张图看清:
run_task 循环(agents.py:1363)
│ execute_task 得到 task_output
▼
_apply_task_guardrail(agents.py:1059)
│ 调 task._process_guardrail(task.py:839)
│ └─ 跑 _guardrail_fn → GuardrailResult
├─ success=False & 还有额度 → retry_count++,返回 should_retry=True
│ → run_task continue,重跑
├─ success=False & 额度用尽 → raise
└─ success=True → (可选)用改写后的产出替换 result → 放行
重点看: guardrail 的重试计数 retry_count / 上限 max_retries 挂在 Task 上(task/task.py:220-221),但驱动重跑的循环在 AgentTeam.run_task。校验逻辑和重试机制是分开的两层。
8. Agent 间交接:Handoff
前面讲的都是"编排器从上帝视角调度任务"。Handoff 是另一个维度的协作:一个 agent 在对话中途,自己决定"这活该转给专职的人"——就像客服把你转接到技术组。
两种触发方式:
- LLM 驱动 —— handoff 被包装成一个工具
transfer_to_<agent_name>(agent/handoff.py:398-402)暴露给 LLM。模型觉得该转就"调用"这个工具,to_tool_function(agent/handoff.py:872)执行真正的交接。 - 编程驱动 —— 直接调
execute_programmatic(agent/handoff.py:630)/execute_async(agent/handoff.py:738)。
交接时传多少上下文,由策略定(ContextPolicy,agent/handoff.py:44-49):
| 策略 | 传什么 | 场景 |
|---|---|---|
NONE | 什么都不传 | 全新任务 |
LAST_N | 最近 N 条(+可选系统消息) | 只需近期上下文 |
SUMMARY(默认) | 系统消息 + 最近 3 条 | 安全默认 |
FULL | 完整历史 | 需要全部背景 |
两个安全设计是精华:
① 工具边界默认收紧(intersect)。 交接时目标 agent 拿到的工具,默认是"源和目标工具的交集"(HandoffToolPolicy.mode="intersect",agent/handoff.py:66)。_compute_effective_tools(agent/handoff.py:413)实现:intersect 模式下,目标只保留"它自己有、且源也有、且不在 blocked_tools"的工具(agent/handoff.py:440-457)。这是"安全优先"——防止交接成了提权。想要老行为(目标保留全部工具)得显式选 passthrough。
② 环路与深度防护。 用 contextvars 记一条"交接链"(_handoff_chain_var,agent/handoff.py:174-176;特意不用 threading.local,避免同线程并发协程互踩)。每次交接前 _check_safety(agent/handoff.py:512):目标已在链上 → 抛 HandoffCycleError;深度超 max_depth(默认 10)→ 抛 HandoffDepthError。防止 A→B→A 无限转接。
结构化交接(TypedHandoff): 普通 handoff 是拼字符串 prompt;TypedHandoff(agent/handoff.py:1289)让发送方声明一个 Pydantic schema,在边界上 _validate_payload 校验(agent/handoff.py:1377),校验过才把规范 JSON(而非 str() 拼接)喂给目标 agent。不合格抛 HandoffValidationError。这让 agent 间传的是"结构化数据"而不是"一段可能被误解析的文本"。
parallel_handoffs(agent/handoff.py:1182)则用信号量控并发,一次把任务扇出给多个目标 agent。
9. 巧妙之处(可借鉴的技术)
-
调度与执行解耦成 generator 协议。
Process只yield"下一个 task_id",AgentTeam只管"跑一条"。换调度策略 = 换一个 generator。三种策略共用同一套执行/重试/guardrail 机制。(agents/agents.py:1899-1915) -
经理 LLM 的双层容错。 让 LLM 出结构化决策不可靠,于是"结构化输出→JSON 模式"双路 fallback(
process/process.py:526)+ "选了非法 task 就回灌合法列表重问"(process/process.py:1715-1741)。把不可靠的 LLM 调用包成可用的调度器。 -
验证反馈闭环。 任务被判失败时,把"错在哪、被拒产出"打包成
validation_feedback挂到重试任务上(process/process.py:925-937),下轮拼进 prompt。重试带记忆,不是盲目重来。 -
交接安全默认收紧。 工具默认取交集、默认摘要上下文、默认开环路检测(
agent/handoff.py:66/94/107)。把"安全"设成缺省而非可选项。 -
双重检查锁复用。 任务记忆懒加载(
task/task.py:567)、异步状态锁(process/process.py:69-75)、handoff 信号量(agent/handoff.py:378-382,现改实例级并绑定事件循环)都用同一套"无锁判空→拿锁→再判"模式防并发重复初始化。