Skip to content

第 21 章 · Workflow 进阶:并行/分支/循环

本章目标:掌握 ParallelConditionLoopRouter 四大控制流原语,能组合它们搭出带质量闭环的内容生产流水线。

21.1 Parallel:并行执行独立步骤

Parallel 块内的步骤并发执行,输出按配置顺序聚合后传给下一步:

python
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 反向分支:

python
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):

python
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 中挑选:

python
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 对象:

python
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) 实现可测试的确定性分支;
  • Loopend_condition + max_iterations 收敛,forward_iteration_output=True 让迭代基于上一轮结果累积;
  • Router.selector 返回要执行的 Step 列表,适合多路择一;
  • 步骤级 output_schema 让结构化数据在整条流水线中流转。

🧪 随堂测验

点击你认为正确的选项。答错时会展示正确答案与原因解析。

1. Parallel 块中多个步骤的输出传给下一步时,顺序是?

2. Loop 在什么情况下停止迭代?

3. 想让 Loop 的下一轮迭代基于上一轮的输出继续加工,应设置?

4. 「根据主题在三种研究策略中选一条执行路径」,最合适的原语是?

🛠️ 动手实践

  1. 给第 20 章的写作流水线加一个 Parallel 调研阶段(两个不同工具的研究员),验证聚合输出的拼接顺序。
  2. 实现"字数不足就重写"的 Loop:end_condition 检查草稿是否超过 200 字,max_iterations=3,观察 forward 开关前后行为差异。
  3. 用 Router 实现"技术问题→HN 研究员 / 生活问题→通用研究员",分别提交两类问题验证路由结果。

流水线跑通之后,下一章把它变成一个 HTTP 服务:AgentOS。