第 21 章 · Workflow 进阶:并行/分支/循环
本章目标:掌握
Parallel、Condition、Loop、Router四大控制流原语,能组合它们搭出带质量闭环的内容生产流水线。
21.1 Parallel:并行执行独立步骤
Parallel 块内的步骤并发执行,输出按配置顺序聚合后传给下一步:
import os
from agno.agent import Agent
from agno.models.openai import OpenAIChat
from agno.workflow import Parallel, Step, StepInput, StepOutput, Workflow
model = OpenAIChat(
id="deepseek-chat",
api_key=os.getenv("DEEPSEEK_API_KEY"),
base_url="https://api.deepseek.com/v1",
)
hackernews_researcher = Agent(name="HN Researcher", role="从 HackerNews 找技术热点", model=model)
web_researcher = Agent(name="Web Researcher", role="搜索网页与官方资料", model=model)
paper_researcher = Agent(name="Paper Researcher", role="检索学术论文观点", model=model)
workflow = Workflow(
name="parallel-research",
steps=[
Parallel(
Step(name="hn", executor=hackernews_researcher),
Step(name="web", executor=web_researcher),
Step(name="papers", executor=paper_researcher),
name="research",
),
Step(name="synthesis", executor=writer), # 拿到三路聚合结果
],
)三个研究方向互不依赖,并行执行能把总耗时压缩到最慢一路的水平。注意:如果并行的函数步骤要写 session state,各分支应写不同的 key,避免竞态。
21.2 Condition:条件分支
Condition 用一个 evaluator 函数决定走哪条路,支持 else_steps 反向分支:
from agno.workflow import Condition, Step, StepInput, StepOutput, Workflow
def is_technical_issue(step_input: StepInput) -> bool:
text = str(step_input.input or "").lower()
return any(kw in text for kw in ["报错", "bug", "崩溃", "超时"])
def diagnose(step_input: StepInput) -> StepOutput:
return StepOutput(content=f"技术诊断: {step_input.input}")
def general_support(step_input: StepInput) -> StepOutput:
return StepOutput(content=f"常规客服回复: {step_input.input}")
workflow = Workflow(
name="support-router",
steps=[
Condition(
name="triage",
evaluator=is_technical_issue,
steps=[Step(name="diagnose", executor=diagnose)], # True 分支
else_steps=[Step(name="general", executor=general_support)], # False 分支
),
Step(name="follow_up", executor=log_and_reply),
],
)evaluator 可以是同步/异步 Python 函数、布尔值或 CEL 表达式字符串。与 Team 的"模型自己决定派谁"相比,这里的分支逻辑是你写的代码——可测试、可复现。
21.3 Loop:迭代打磨直到达标
Loop 重复执行内部步骤,直到 end_condition 返回 True 或达到 max_iterations(默认 3):
from agno.workflow import Loop, Step, StepInput, StepOutput, Workflow
def draft_section(step_input: StepInput) -> StepOutput:
prev = step_input.previous_step_content or str(step_input.input)
return StepOutput(content=f"{prev}\n【补充一段更具体的论据】")
def quality_check(outputs: list[StepOutput]) -> bool:
"""所有迭代输出中任一超过 120 字即认为达标"""
return any(len(str(o.content or "")) > 120 for o in outputs)
def final_edit(step_input: StepInput) -> StepOutput:
return StepOutput(content=f"终稿:\n{step_input.previous_step_content}")
workflow = Workflow(
name="quality-loop",
steps=[
Loop(
name="drafting-loop",
steps=[Step(name="expand", executor=draft_section)],
end_condition=quality_check,
max_iterations=3,
forward_iteration_output=True, # 下一轮拿到上一轮的输出作为 previous_step_content
),
Step(name="edit", executor=final_edit),
],
)forward_iteration_output=True 是迭代累积的关键:默认每轮都从原始输入重新开始;开启后,下一轮的 previous_step_content 是上一轮的结果,实现"越写越长、逐轮打磨"。
21.4 Router:动态选择执行路径
Router 的 selector 函数返回要走的一组步骤,从 choices 中挑选:
from typing import List
from agno.workflow import Router, Step
from agno.workflow.types import StepInput
tech_step = Step(name="tech_research", executor=hackernews_researcher)
general_step = Step(name="general_research", executor=web_researcher)
def research_router(step_input: StepInput) -> List[Step]:
topic = (step_input.previous_step_content or step_input.input or "").lower()
if any(kw in topic for kw in ["ai", "编程", "startup", "软件"]):
return [tech_step] # 返回列表,可以返回多个步骤
return [general_step]
workflow = Workflow(
name="router-demo",
steps=[
Router(
name="strategy",
selector=research_router,
choices=[tech_step, general_step],
),
publish_step,
],
)选型口诀:布尔判断用 Condition,多路择一用 Router。两者都能嵌套组合——例如 Router 选出的分支里再放 Loop。
21.5 组合实战:内容生产 Pipeline + 步骤级结构化 IO
把本章原语串成一条真实流水线,并让每个 LLM 步骤用 output_schema 输出 Pydantic 对象:
from pydantic import BaseModel, Field
class Draft(BaseModel):
title: str = Field(description="文章标题")
body: str = Field(description="正文")
word_count: int
class Review(BaseModel):
approved: bool
comments: list[str]
drafter = Agent(name="Drafter", model=model, output_schema=Draft,
instructions=["按 schema 输出初稿"])
review_agent = Agent(name="Reviewer", model=model, output_schema=Review,
instructions=["审核初稿,approved 表示是否通过"])
content_workflow = Workflow(
name="content-pipeline",
steps=[
Parallel(Step(name="collect_hn", executor=hackernews_researcher),
Step(name="collect_web", executor=web_researcher),
name="collect"),
Step(name="draft", executor=drafter),
Loop(
name="revise-loop",
steps=[Step(name="review", executor=review_agent)],
end_condition=lambda outputs: all(
getattr(o.content, "approved", False) for o in outputs
),
max_iterations=3,
forward_iteration_output=True,
),
Step(name="publish", executor=publisher),
],
)每个步骤的输出都是经过 Pydantic 校验的对象,下游步骤拿到的不再是"祈祷格式正确"的纯文本——这是生产级 Workflow 与 demo 的分水岭。
类执行器
需要初始化配置或维护状态时,可以实现 __call__(self, step_input) -> StepOutput 的类作为 executor,实例在多次运行间复用。
21.6 本章小结
Parallel并发跑独立步骤、按序聚合输出;并行分支写共享 state 要分开 key;Condition(evaluator, steps, else_steps)实现可测试的确定性分支;Loop靠end_condition+max_iterations收敛,forward_iteration_output=True让迭代基于上一轮结果累积;Router.selector返回要执行的 Step 列表,适合多路择一;- 步骤级
output_schema让结构化数据在整条流水线中流转。
🧪 随堂测验
点击你认为正确的选项。答错时会展示正确答案与原因解析。
1. Parallel 块中多个步骤的输出传给下一步时,顺序是?
2. Loop 在什么情况下停止迭代?
3. 想让 Loop 的下一轮迭代基于上一轮的输出继续加工,应设置?
4. 「根据主题在三种研究策略中选一条执行路径」,最合适的原语是?
🛠️ 动手实践
- 给第 20 章的写作流水线加一个
Parallel调研阶段(两个不同工具的研究员),验证聚合输出的拼接顺序。 - 实现"字数不足就重写"的 Loop:
end_condition检查草稿是否超过 200 字,max_iterations=3,观察 forward 开关前后行为差异。 - 用 Router 实现"技术问题→HN 研究员 / 生活问题→通用研究员",分别提交两类问题验证路由结果。
流水线跑通之后,下一章把它变成一个 HTTP 服务:AgentOS。