第 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
# 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:
# 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 调用方式:
# 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() 启动:
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 再辨析
| 维度 | Crew | Flow |
|---|---|---|
| 抽象层次 | 角色协作团队 | 流程编排骨架 |
| 控制流 | 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?
🛠️ 动手实践
- 把第 5 章的双角色写作 Crew 改造成 Flow:
@start接收主题 → 内嵌 Crew 产出文章 →@listen统计字数并写入article.md。 - 用 Pydantic 结构化 state 重写本章 CityFlow,要求 state 中显式声明
city与fun_fact字段,并在结束时打印完整 state。 - 写一个含两个
@start方法的 Flow(如同时抓取天气与新闻),观察两者的执行顺序,并把结果都存进 state 后用一个@listen(and_(...))方法汇总(预告下一章的合流技巧)。
单靠返回值和 state 还不够,下一步学习条件路由与持久化:第 16 章 · Flow 进阶:状态管理与持久化。