确定性工作流 AgentFlow 与条件系统
30 秒导读:
AgentFlow是 PraisonAI 的第二条编排范式——开发者把执行顺序写死(用 Python 代码或 YAML),框架照单执行,不需要一个 manager LLM 去决定下一步做什么。它给你六种组合积木(顺序、分支、并行、循环、重试、复用),加一套用正则解析字符串表达式的条件引擎来做路由。
本章讲"确定性 DAG / 循环"这条范式本身。第 04 章的 Task / AgentTeam / 三种 Process 是另一条范式(manager 决策式),本章只在需要对照时提它,不重复它的内容。
1. 这是什么(零基础也能懂)
一 句话定义: AgentFlow 是一条你自己写死顺序的智能体流水线——第一步跑完喂给第二步,中间可以插分支、并行、循环,全部由你在代码里排好,框架只负责忠实执行。
解决什么问题 / 给谁用: 假设你要做一个"写文章 → 编辑 → 发布"的固定流程。你已经知道这三步的先后,不需要让一个 LLM 每次去"想一想现在该干嘛"——那样既慢又不可控。AgentFlow 让你像写普通函数调用链一样,把流程钉死下来。
两条范式的分工(这是理解本章的关键):
| 范式 | 谁决定"下一步做什么" | 适合 |
|---|---|---|
AgentTeam + Process(第 04 章) | 一个 manager LLM 在运行时决策 | 步骤不固定、要智能调度的开放任务 |
AgentFlow(本章) | 开发者在代码/YAML 里写死 | 步骤已知、要可复现、要可控成本的流水线 |
用起来什么样: 最小的顺序流水线——把两个 Agent 依次串起来,run() 一把跑完。
# 示意,非源码(真实 API 见 workflows.py:555 AgentFlow / :998 run)
from praisonaiagents import AgentFlow, Agent
flow = AgentFlow(steps=[
Agent(instructions="写一段关于 AI 的内容"), # 第 1 步
Agent(instructions="把上一步的内容润色"), # 第 2 步:自动收到上一步输出
])
result = flow.run("写 AI") # 或 flow.start(...),两者等价
print(result["output"]) # 最后一步的输出
一句话直觉: 把它当成 Unix 管道 a | b | c——数据从左流到右,每一节是一个智能体或一个函数。只不过这条管道还能长出分支(if)、分叉(parallel)和回 环(loop/repeat)。
2. 顶层全景(它大概怎么转)
2.1 核心部件
| 部件 | 干什么 | 在哪 |
|---|---|---|
AgentFlow | 工作流主体,持有 steps 列表,run() 驱动执行 | workflows.py:555 |
WorkflowContext | 传给每一步的只读上下文(input、上一步输出、变量表) | workflows.py:174 |
StepResult | 每一步返回的结果(output、是否提前终止、要写回的变量) | workflows.py:182 |
| 六种组合原语 | Route/Parallel/Loop/Repeat/If/Include——控制流积木 | workflows.py:197–551 |
| 条件引擎 | evaluate_condition() 把 "{{score}} > 80" 求值成布尔 | conditions/evaluator.py:129 |
YAMLWorkflowParser | 把 YAML 文件解析成一个 AgentFlow | workflows/yaml_parser.py:21 |
2.2 主循环:一个"边走边认类型"的调度器
AgentFlow.run()(workflows.py:998)的心脏是一个 while i < len(self.steps) 循环(workflows.py:1092)。它逐个取出 steps 里的元素,先看它是不是某种组合原语,是就交给对应的 _execute_* 处理器;否则当成普通单步(Agent / 函数 / Task)执行。
怎么读下图:从上往下是主循环每一轮的判断顺序,命中一种就分派、然后 i += 1 进入下一轮。
run(input) ──► while i < len(steps): 取 step = steps[i]
│
├─ 是 Route? ──► _execute_route (按关键词路由)
├─ 是 Parallel? ──► _execute_parallel (线程池并发)
├─ 是 Loop? ──► _execute_loop (遍历列表/CSV)
├─ 是 Repeat? ──► _execute_repeat (重复到满足 until)
├─ 是 Include? ──► _execute_include (嵌入另一个 recipe)
├─ 是 If? ──► _execute_if (表达式真→then 假→else)
│
└─ 都不是 ──► 普通单步:Agent.chat() / 函数 handler / 临时 Agent
└─ 输出写入 previous_output,并存进变量表
依据:分派判断在
workflows.py:1096-1160;普通单步执行在workflows.py:1162-1483。
三条贯穿全程的暗线(后面各节展开):
- 输出即输入。 每步的
output存进previous_output,下一步默认自动收到它(workflows.py:1448、变量替换见_substitute_action_variablesworkflows.py:107)。 - 变量表
all_variables一路累积。 每步结果按f"{step.name}_output"或自定义output_variable写回(workflows.py:1460),供后续条件/模板引用。 - 确定性。 走哪条分支、循环几次,只取决于变量值和写定的结构,没有 manager LLM 在中间拍板。
3. 组合原语(六种控制流积木)
这是本章最核心的部分。AgentFlow 用六个积木拼出任意确定性控制流。每个积木都有两种写法:小写便捷函数(route(...))和大写数据类(Route(...))——便捷函数只是薄封装,最终都产出同一个数据类实例。
| 便捷函数 | 数据类 | 作用 | 类比 |
|---|---|---|---|
route() | Route | 按上一步输出里的关键词跳到不同分支 | switch/case |
parallel() | Parallel | 多步并发跑完再汇总 | fork/join |
loop() | Loop | 对列表 / CSV / 文件逐项执行 | for-each |
repeat() | Repeat | 重复同一步直到条件满足 | do-while |
when() / if_() | If | 求值表达式,真走 then 假走 else | if/else |
include() | Include | 把另一个 recipe / workflow 当一步嵌入 | 函数调用 |
依据:便捷函数
route/parallel/loop/repeat/include在workflows.py:361-456;when/if_在workflows.py:506-551;对应数据类在workflows.py:197/219/250/332/410/461。__init__.py:20-35把它们全部导出。
3.1 Route —— 按关键词分支
要解决的小问题: 上一步(通常是个"决策"智能体)输出了一段话,里面含 "approve" 或 "reject",我想据此走不同后续。
思路: 不做复杂求值,直接在上一步输出文本里搜关键词——但用的是词边界匹配(\bkey\b),避免 "approved" 里的子串误命中。
# 示意,非源码
from praisonaiagents.workflows import route
route({
"approve": [publish_agent], # 输出含 "approve" → 走这条
"reject": [revise_agent],
"default": [fallback_agent], # 都不含 → 兜底
})
真实实现: _execute_route() 遍历 route 键,用 re.search(r'\b'+key+r'\b', prev_lower) 命中即停,没命中走 default(workflows.py:2264-2274)。
关键细节: Route 匹配的是上一步的输出文本,不是变量表里的值——这点和下面的 If 正好相反,别混。
3.2 Parallel —— 并发分叉再汇总
要解决的小问题: 三个互不依赖的子任务,串行跑太慢。
思路: 用 ThreadPoolExecutor 并发,全部跑完把输出用 \n---\n 拼起来(workflows.py:2442),并存进 parallel_outputs 变量。
两个要注意的设计:
- 默认限流。 未指定
max_workers时,并发数 =min(DEFAULT_MAX_PARALLEL_WORKERS, 分支数),而DEFAULT_MAX_PARALLEL_WORKERS = 3(workflows.py:42、:2399-2400)——刻意压着,防止 LLM 后端被打到限流。 - 三种失败策略(
on_failure,在Parallel.__init__校验,非法值直接抛错workflows.py:240-245):
| 值 | 语义 |
|---|---|
partial_ok(默认) | 某分支失败也继续,把错误当该分支输出 |
fail_fast | 首个失败即取消其余分支并抛 WorkflowStepError |
fail_all | 等所有分支跑完,只要有失败就抛 |
依据:失败分流在
workflows.py:2424-2439。
3.3 Loop —— 逐项遍历
Loop 对一个列表变量(over="items")、一个 CSV(from_csv)或文本文件(from_file)逐项执行;可单步也可多步(steps=[...]),可串行也可 parallel=True 并发(workflows.py:250-330)。构造时就校验"不能同时给 step 和 steps、也不能都不给"(workflows.py:313-320)。
3.4 Repeat —— 重复到满足条件(evaluator-optimizer)
要解决的小问题: "生成 → 自检 → 不够好就再生成",最多试 N 次。
思路: 反复跑同一步,每轮后调用 until 回调判断是否收敛;until 是个Python 可调用对象(接收 WorkflowContext 返回 bool),到达 max_iterations(默认 10)无条件停。
# 示意,非源码
from praisonaiagents.workflows import repeat
repeat(generator,
until=lambda ctx: "done" in ctx.previous_result.lower(),
max_iterations=5)
真实实现: _execute_repeat() 的 for iteration in range(max_iterations) 循环,每轮跑完构造 WorkflowContext 再调 until(workflows.py:2750-2772)。注意:这里 until 是代码回调,不是字符串表达式——和下面 If 的字符串条件不是一套东西。
3.5 If / when —— 表达式真假分支
要解决的小问题: "如果分数 > 80 就批准,否则打回"——这次判断依据是变量表里的值,不是文本关键词。
思路: when() 是 if_() 的首选别名(两者完全等价,都造 If 对象,workflows.py:506/533)。它拿一个字符串条件 "{{score}} > 80",交给条件引擎求值成布尔,真走 then_steps 假走 else_steps。
# 示意,非源码
from praisonaiagents.workflows import when
when(condition="{{score}} > 80",
then_steps=[approve_agent],
else_steps=[reject_agent])
真实实现: _execute_if() 先 _evaluate_condition(...) 拿布尔,再选 分支执行(workflows.py:2810-2817)。条件引擎是本章第 4 节的主角。
3.6 Include —— 模块化复用
include("wordpress-publisher") 或 include(workflow=other_flow) 把另一个 recipe / workflow 当一步嵌进来,实现模块化组合;构造时强制"recipe 和 workflow 至少给一个"(workflows.py:438-443)。执行时带环检测:同一执行链里重复 include 同名 recipe 会被拦下报 "Circular include detected"(workflows.py:2894-2901)。
3.7 嵌套与深度上限
原语可以互相嵌套(if 里放 parallel,loop 里放 route……)。统一入口 _execute_single_step_internal()(workflows.py:2045)在递归进入嵌套原语时把 depth+1,一旦 depth > MAX_NESTING_DEPTH(=5,workflows.py:459)就抛错,防止无限递归爆栈(workflows.py:2074-2078)。
_execute_single_step_internal(step, depth)
│ depth > 5 ? ──► ValueError("Maximum nesting depth exceeded")
├─ Loop ─► _execute_loop(..., depth+1)
├─ Parallel ─► _execute_parallel(..., depth+1)
├─ Route ─► _execute_route(..., depth+1)
├─ Repeat ─► _execute_repeat(..., depth+1)
├─ If ─► _execute_if(..., depth+1)
└─ 普通步 ─► normalize → Agent/handler/临时Agent
依据:嵌套分派
workflows.py:2081-2148。
4. 条件求值系统(字符串表达式怎么被安全求值)
这是与"确定性 DAG"并列的第二个引擎。praisonaiagents/conditions/ 独立成模块,被 AgentFlow(字符串条件)和 AgentTeam(字典路由)共用,做 DRY 复用。