Skip to content

第 15 章 · Flow 入门:事件驱动工作流

本章目标:理解 CrewAI Flow 的事件驱动模型,掌握 Flow 类与 @start / @listen 装饰器、方法间经返回值和 state 传递数据,学会在 Flow 中内嵌 Crew,并能依据场景在 Flow 与 Crew 之间做正确选型。

15.1 Crew 的天花板与 Flow 的定位

Crew 擅长"一组角色协作完成一个目标",但它有两个天然短板:

  • 控制流弱:顺序/分层两种流程难以表达"如果 A 失败就走 B"、"先并行跑 X 和 Y 再合流"这类编程逻辑;
  • 状态管理散:跨 Crew 共享数据只能靠任务输出层层传递。

Flow 是 CrewAI 对这个问题的答案——结构化的、事件驱动的工作流:你用普通的 Python 方法定义步骤,用装饰器声明"谁触发谁",数据通过 state 与返回值流动。官方总结的三大价值:简化多 Crew/任务编排、内置状态管理、灵活的控制流(条件、分支)。

15.2 第一个 Flow:@start 与 @listen

python
# first_flow.py —— 最小 Flow:生成城市 → 生成趣闻
import os
from crewai import LLM
from crewai.flow.flow import Flow, listen, start

llm = LLM(
    model="openai/deepseek-chat",
    base_url="https://api.deepseek.com/v1",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
)

class CityFlow(Flow):

    @start()                      # 标记流程入口
    def generate_city(self):
        print("Flow 启动")
        # 每个 Flow 的 state 会自动获得唯一 ID
        print(f"State ID: {self.state['id']}")
        result = llm.call("随便返回世界上一个城市名。")
        self.state["city"] = result          # 写入共享状态
        return result                        # 返回值也会传给监听者

    @listen(generate_city)        # 监听上一步:它完成后才执行
    def generate_fun_fact(self, random_city):
        fun_fact = llm.call(f"讲一个关于 {random_city} 的有趣事实。")
        self.state["fun_fact"] = fun_fact
        return fun_fact


flow = CityFlow()
flow.plot("city_flow_plot")       # 生成可视化 HTML(可选)
result = flow.kickoff()           # 执行,返回最后完成方法的输出
print(f"生成的趣闻: {result}")

要点拆解:

  • @start() 标记入口方法;Flow 启动时所有满足条件的 @start 方法都会执行(多个时往往并行);
  • @listen(xxx) 声明"当 xxx 完成后执行我";被监听方法的返回值会作为参数传入;
  • kickoff() 返回最后一个完成的方法的输出——上面就是 generate_fun_fact 的返回值;
  • self.state 是所有方法共享的状态容器,未定义结构时是字典,且自动带一个 id 键(UUID)。

@listen 还支持字符串形式:@listen("generate_city") 等价于传方法对象,适合类中前向引用的场景。

15.3 结构化 state:TypedDict / Pydantic

字典式 state 灵活但容易写错键名。给 Flow 加泛型参数即可升级为结构化状态,推荐 Pydantic BaseModel

python
# state_flow.py —— Pydantic 结构化状态
from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, start

class PoemState(BaseModel):
    sentence_count: int = 1     # 注意:id 字段会被自动添加并维护
    poem: str = ""

class PoemFlow(Flow[PoemState]):

    @start()
    def prepare(self):
        self.state.sentence_count = 4      # 属性访问 + IDE 自动补全
        self.state.poem = ""

    @listen(prepare)
    def write_poem(self):
        text = llm.call(f"写一首 {self.state.sentence_count} 句的短诗。")
        self.state.poem = text
        return text

flow = PoemFlow()
final = flow.kickoff()
print(flow.state)      # 查看最终状态: sentence_count=4 poem='...'

选择建议(官方口径):state 结构简单多变、追求原型速度 → 非结构化字典;需要类型安全、校验和 IDE 提示 → Pydantic 结构化状态。生产项目几乎都应该用后者。

15.4 在 Flow 方法里调用 LLM 与内嵌 Crew

Flow 方法里可以自由混用三种 AI 调用方式:

python
# mixed_flow.py —— 三种调用方式并存
import os
from pydantic import BaseModel
from crewai import Agent, Task, Crew, Process, LLM
from crewai.flow.flow import Flow, listen, start

llm = LLM(
    model="openai/deepseek-chat",
    base_url="https://api.deepseek.com/v1",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
)

class ReportState(BaseModel):
    topic: str = ""
    outline: str = ""
    report: str = ""

class ReportFlow(Flow[ReportState]):

    @start()
    def set_topic(self):
        self.state.topic = "AI Agent 在客服系统的落地"

    @listen(set_topic)
    def make_outline(self):
        # 方式一:直接调 LLM —— 轻量单步任务
        self.state.outline = llm.call(
            f"为《{self.state.topic}》拟一个三段式大纲,只输出大纲。"
        )
        return self.state.outline

    @listen(make_outline)
    def write_report(self):
        # 方式二:单个 Agent 直接 kickoff —— 不需要团队协作时更省事
        writer = Agent(role="撰稿人", goal="按大纲成文",
                       backstory="技术写作专家", llm=llm)
        t = Task(description=f"按大纲撰写报告:\n{self.state.outline}",
                 expected_output="800 字左右的 Markdown 报告。",
                 agent=writer)
        single = Crew(agents=[writer], tasks=[t],
                      process=Process.sequential)
        out = single.kickoff()
        self.state.report = out.raw
        return out.raw


if __name__ == "__main__":
    result = ReportFlow().kickoff()
    print(result)

第三种方式是把完整的多角色 Crew 封装成一个函数供 Flow 步骤调用(下一章会用脚手架项目的形态展开)。判断标准:单步能搞定就 LLM.call 或单 Agent;需要多角色协作、工具分工才上 Crew。

15.5 异步方法与 kickoff_async

Flow 方法可以是普通函数,也可以是 async def;对应地用 await flow.kickoff_async() 启动:

python
import asyncio
from crewai.flow.flow import Flow, listen, start

class AsyncFlow(Flow):

    @start()
    async def fetch_data(self):           # 异步方法:可做 IO 密集操作
        await asyncio.sleep(0.5)          # 模拟异步请求
        self.state["data"] = "原始数据"
        return "ok"

    @listen(fetch_data)
    def process(self, status):            # 同步方法也可以监听异步方法
        return f"{status}: {self.state['data']}"

asyncio.run(AsyncFlow().kickoff_async())

Web 服务集成(FastAPI 接口触发工作流)时,kickoff_async 能避免阻塞事件循环,这是生产部署的标准姿势。

15.6 Flow vs Crew 再辨析

维度CrewFlow
抽象层次角色协作团队流程编排骨架
控制流sequential / hierarchical任意 Python 逻辑 + 条件路由
状态任务输出链显式 state(可持久化)
适用"一件事交给一个团队""一条业务流水线,含多个团队/步骤"

经验法则:先用最简单的方案——单次问答不需要任何框架概念;多角色协作一次成型用 Crew;出现条件分支、并行汇合、断点恢复、跨阶段共享状态中的任意一种需求,就该把 Crew 装进 Flow 里了。

15.7 本章小结

  • Flow 用 @start() 声明入口、@listen() 声明依赖,构成事件驱动的有向图;kickoff() 返回最后完成方法的输出;
  • 数据两条通路:方法返回值 → 监听者入参;self.state 全局共享(字典或 Flow[PydanticModel] 泛型);
  • Flow 方法内可混合 LLM.call、单 Agent、完整 Crew 三种粒度的 AI 调用;
  • 异步方法配 kickoff_async(),服务化部署必备;
  • 选型:协作一次成型用 Crew,有分支/并行/状态/恢复需求用 Flow 包住 Crew。

🧪 随堂测验

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

1. flow.kickoff() 的返回值是什么?

2. @listen(generate_city) 的监听方法如何拿到上一步的数据?

3. 想要 state 具备字段校验和 IDE 自动补全,应该怎么做?

4. 以下哪种场景最适合直接用 Crew 而不必引入 Flow?

🛠️ 动手实践

  1. 把第 5 章的双角色写作 Crew 改造成 Flow:@start 接收主题 → 内嵌 Crew 产出文章 → @listen 统计字数并写入 article.md
  2. 用 Pydantic 结构化 state 重写本章 CityFlow,要求 state 中显式声明 cityfun_fact 字段,并在结束时打印完整 state。
  3. 写一个含两个 @start 方法的 Flow(如同时抓取天气与新闻),观察两者的执行顺序,并把结果都存进 state 后用一个 @listen(and_(...)) 方法汇总(预告下一章的合流技巧)。

单靠返回值和 state 还不够,下一步学习条件路由与持久化:第 16 章 · Flow 进阶:状态管理与持久化