Flow 是什么:给 Agent 世界装一套"事件驱动的流程引擎"
前 40 天我们把 Agent / Task / Crew 拆到了底:Crew 会按 sequential 或 hierarchical 把一串任务跑完。但真实业务不是"一条直线"——有分支、有循环、有"等 A 和 B 都好了再做 C"。这类确定性流程正是 Flow 要解决的。今天先鸟瞰 crewai/flow/ 这个包:它被拆成 DSL(写法)/ Definition(契约)/ Runtime(引擎) 三层,公开的 Flow 类怎么拼出来,运行时骨架长啥样,以及一句 kickoff() 是怎么把方法按事件串起来跑的。
Crew 是"把一队人按顺序排好班让他们干活",那 Flow 就是一张流程图 + 一个自动派单的调度台。你在流程图上画好"哪个环节做完了,就通知下一个环节开工"(@start / @listen),调度台(Runtime)盯着每个环节完成的事件,一旦某个环节发出"我做完了"的信号,它就去查"谁在等这个信号",把等到的环节全叫起来跑。整个 Flow 的运转就是「方法完成 → 发事件 → 查谁在监听 → 触发下一批方法」的循环,而 Crew 只是流程图里可以被塞进某个环节的一块"积木"。痛点:Crew 为什么编排不了"带分支的流程"?
@start/@listen/@router 装饰的方法是图上的一个节点,节点之间靠"某方法完成"这个事件连边。kickoff() 启动后,引擎不停地做「执行方法 → 该方法完成的事件广播出去 → 找到监听这个事件的方法 → 执行它们」,直到没有方法可触发。控制流(分支/循环/汇合)由"谁监听谁"表达,业务逻辑就是普通 Python,LLM 只在你显式调用 Crew/Agent 时才介入。flow/ 目录:一件事拆成三层
打开 crewai/flow/flow.py,开头的模块 docstring 就把整个包的分层说清楚了(flow/flow.py:1):
# flow/flow.py:1(模块 docstring 摘要)
"""Backwards-compatible re-export surface for the Flow framework.
The implementation now lives in three modules, split by concern:
- ``crewai.flow.dsl`` -- 写法层:@start / @listen / @router、or_ / and_
- ``crewai.flow.flow_definition`` -- 契约层:可序列化的 Flow Definition
- ``crewai.flow.runtime`` -- 引擎层:Flow 执行引擎与状态
- ``crewai.experimental.conversational_mixin`` -- 实验性的对话式扩展
"""
dsl(写法层)你写 Flow 类时用的那些装饰器都在这。它的产物是一份 FlowDefinition——把"你画的图"翻译成一个纯数据结构。flow_definition(契约层)用 Pydantic 定义"一个 Flow 长什么样":有哪些方法、谁监听谁、状态是什么。可序列化意味着 Flow 能存成 YAML/JSON,不写 Python 也能声明一个 Flow。runtime(引擎层)真正的执行引擎:kickoff、事件广播、状态管理都在这个 3600 行的大文件里。flow.py 本身只是个兼容门面:把三层的东西重新导出,保留历史上 from crewai.flow.flow import Flow 这个老路径。Flow.from_declaration)。中间隔一层契约(Definition),写法层和引擎层就彻底解耦:换写法不影响引擎,换引擎不影响写法。代价是多了一次"翻译"(类 → Definition),换来的是可序列化、可可视化、可声明式三大能力。你 import 的那个 Flow,是怎么拼出来的?
你平时写 class MyFlow(Flow) 用的 Flow,其实是三个来源组合出来的(flow/flow.py:33):
# flow/flow.py:20
from crewai.experimental.conversational_mixin import _ConversationalMixin
from crewai.flow.dsl import and_, listen, or_, router, start
from crewai.flow.runtime import (
_INITIAL_STATE_CLASS_MARKER,
Flow as RuntimeFlow, # ← 引擎层的 Flow 才是真身
FlowMeta,
FlowState,
)
T = TypeVar("T", bound=dict[str, Any] | BaseModel)
# flow/flow.py:33
class Flow(_ConversationalMixin, RuntimeFlow[T]):
"""Public Flow class with experimental conversational extension behavior."""
RuntimeFlow[T]公开 Flow 的真正基类是引擎层的 Flow。业务能力(kickoff、状态、事件)全来自它。_ConversationalMixin再叠一层"对话式"混入(D47 讲)——让 Flow 能当聊天机器人用。放在前面,方法解析优先级更高。Generic[T]T 被约束为 dict 或 BaseModel——这就是 Flow 的状态类型。Flow[MyState] 让状态带上强类型。start/listen/router/and_/or_装饰器和组合器也在这里重新导出,所以 from crewai.flow.flow import start 能用。RuntimeFlow,把可选的、实验性的对话能力做成一层可插拔的 Mixin,将来要去掉对话能力,改一行基类列表就行,不动引擎。运行时 Flow 的骨架:一个"带状态的 Pydantic 模型"
引擎层的 Flow(runtime/__init__.py:428),注意它继承的是 BaseModel:
# runtime/__init__.py:428
class Flow(BaseModel, Generic[T], metaclass=FlowMeta):
"""Base class for all flows.
type parameter T must be either dict[str, Any] or a subclass of BaseModel."""
model_config = ConfigDict(
arbitrary_types_allowed=True,
ignored_types=(StartMethod, ListenMethod, RouterMethod), # 装饰器包装不当字段
revalidate_instances="never",
)
__hash__ = object.__hash__
_flow_definition: ClassVar[FlowDefinition | None] = None # 懒构建、按类缓存
entity_type: Literal["flow"] = "flow"
_state: Any = PrivateAttr(default=None) # ← 运行时状态存这
继承 BaseModelFlow 本身是一个 Pydantic 模型!这样 Flow 的配置字段(stream、max_method_calls、persistence…)能被 Pydantic 校验,也能 from_declaration 从 dict 反序列化出来。ignored_types=(StartMethod, …)★关键:你的 @start 方法被包装成 StartMethod 对象。告诉 Pydantic"这些不是数据字段,别当模型字段处理",否则会报错。_flow_definition ClassVar解析出来的静态契约按类缓存——同一个 Flow 类只解析一次图结构。_state PrivateAttr运行时的状态对象(你的 self.state)存在私有属性里,不是模型字段,避免被序列化/校验干扰。契约是懒构建的——第一次访问才解析(runtime/__init__.py:469):
# runtime/__init__.py:469
@classmethod
def flow_definition(cls) -> FlowDefinition:
"""Return the static Flow Definition built from this Flow class."""
flow_definition = cls.__dict__.get("_flow_definition")
if flow_definition is None:
flow_definition = build_flow_definition(cls) # 扫描类里的装饰器方法
cls._flow_definition = flow_definition
return flow_definition
FlowMeta.__new__ 里有一句注释点明了原因(runtime/__init__.py:422):"The static FlowDefinition is built lazily … to avoid AST parsing and diagnostic logging on every import." 朴素做法是在元类里一建类就解析图。但解析要扫描方法、读类型注解、可能还打诊断日志——如果每次 import 你的模块都跑一遍,启动会变慢、日志会刷屏。懒构建把这份开销推迟到真正要跑或要画图时,且用 ClassVar 缓存只算一次。"按需 + 缓存"是解析类元信息的标准优化。FlowMeta 元类:替你把"裸赋值"变成合法字段
some_config = 3 这种没有类型注解的裸赋值,或者放一些普通方法。纯 Pydantic 会因为"字段没注解"报错。谁来兜底?答案是自定义元类 FlowMeta(runtime/__init__.py:375),它在类被创建时改写注解:
# runtime/__init__.py:405
for attr_name, attr_value in list(namespace.items()):
if attr_name in annotations or attr_name.startswith("_"):
continue
if attr_name in parent_fields:
annotations[attr_name] = Any
...
continue
if callable(attr_value) or isinstance(
attr_value, (*_skip_types, FlowMethod) # 方法/装饰器包装 → 不当字段
):
continue
annotations[attr_name] = ClassVar[type(attr_value)] # 裸赋值 → 当类变量
namespace["__annotations__"] = annotations
startswith("_") 跳过私有属性不管,交给 Pydantic 的 PrivateAttr 机制。callable / FlowMethod 跳过普通方法、被 @start 等包成 FlowMethod 的方法,都不是数据字段,直接放过。其余裸赋值 → ClassVar★没注解的裸赋值(如 foo = 1)被自动标成 ClassVar,于是 Pydantic 不再把它当"实例字段"要求你注解。parent_fields → Any和父模型字段重名的,标成 Any 以免类型冲突。这是为了让 Flow 能安全地被继承/扩展。kickoff:一句启动,事件驱动怎么转起来
同步入口 kickoff 其实是异步引擎的薄包装(runtime/__init__.py:1920):
# runtime/__init__.py:1966
async def _run_flow() -> Any:
return await self.kickoff_async(inputs, input_files, ...)
runtime_scope = crewai_event_bus._enter_runtime_scope()
try:
try:
asyncio.get_running_loop() # 已在事件循环里?
ctx = contextvars.copy_context()
with ThreadPoolExecutor(max_workers=1) as pool:
return pool.submit(ctx.run, asyncio.run, _run_flow()).result()
except RuntimeError: # 没有运行中的循环
return asyncio.run(_run_flow()) # 直接开一个跑
finally:
crewai_event_bus._exit_runtime_scope(runtime_scope)
本质是包 kickoff_async所有状态初始化、事件广播都在异步版里。同步版只负责"找个事件循环把它跑完"。get_running_loop() 探测★如果调用方本来就在异步环境里(比如 Jupyter、别的 async 代码),不能再 asyncio.run(会报"loop already running"),于是另起一个线程跑,避开冲突。except RuntimeError没有运行中的循环 → 说明是普通脚本,asyncio.run 一把梭。_enter/_exit_runtime_scope进出时给事件总线开/关一个作用域——保证这次 kickoff 发的事件被正确归拢。而 kickoff_async 的核心动作,是执行所有 start 方法、随后让事件链自动扩散(runtime/__init__.py:2442 _execute_start_method):
# runtime/__init__.py:2473(精简)
result, finished_event_id = await self._execute_method(
start_method_name, enhanced_method) # ① 跑一个 @start 方法
...
await self._execute_listeners( # ② 广播"它完成了",触发监听者
start_method_name, result, finished_event_id)
@start def fetch(self): return "raw",@listen(fetch) def clean(self, data): return data.upper()。kickoff → 执行
fetch 得 "raw" → 广播"fetch 完成" → 发现 clean 监听 fetch → 执行 clean("raw") 得 "RAW" → 广播"clean 完成" → 没人监听 clean → 结束,返回 "RAW"。你从头到尾没写一句 while,流程却自己转完了。Flow 与 Crew:不是二选一,是"编排"与"被编排"
看 crewai/flow/__init__.py 导出的公开面,就能感受到 Flow 是个独立的编排层(flow/__init__.py:25):
# flow/__init__.py:25(节选 __all__)
__all__ = [
"Flow", "start", "listen", "router", "and_", "or_", # 图的写法
"FlowStructure", "build_flow_structure", "visualize_flow_structure", # 可视化
"persist", # 状态持久化(D45)
"Expression", # CEL 表达式(D46)
"ConversationalConfig", "ChatState", # 对话式(D47)
"InputProvider", "HumanFeedbackProvider", ... # 人在环 / 输入
]
result = MyCrew().kickoff(inputs=...),把 Crew 当普通函数用。反过来 Crew 里塞不进 Flow——因为 Crew 没有"事件监听"这套控制流能力。| 维度 | Crew | Flow |
|---|---|---|
| 核心抽象 | Agent + Task 列表 + Process | 方法图 + 事件监听 |
| 控制流 | 顺序 / 层级(写死两种) | 分支 / 循环 / 汇合(你画) |
| 谁决策 | LLM 驱动为主 | 你的 Python 代码驱动,LLM 按需介入 |
| 状态 | 任务间 context 传递 | 显式 self.state(可持久化) |
| 定位 | 被编排的"积木" | 顶层"编排者" |
👶 小白:那我到底该用 Crew 还是 Flow?
👨🏫 老师:需求能用"一队 Agent 顺次/层级把任务做完"描述,就用 Crew,简单直接。一旦出现if/else 分支、要循环重试、要等多路汇合、要精确控制每一步,就上 Flow,并在需要"智能"的那一步里调用 Crew。大项目常见结构是:一个 Flow 编排全局,里面几个节点各自 kickoff 一个 Crew。D48 会专门讲选型。
边界 + 今日小结
_create_initial_state 最后有一句兜底(runtime/__init__.py:1609):raise TypeError("Initial state must be dict or BaseModel, got ...")。类型参数 T 在类定义处就被约束为 dict[str, Any] | BaseModel。如果你想拿一个 @dataclass 或普通对象当状态,会在这里炸。为什么这么严?因为状态要能被序列化持久化(D45)、要能被 CEL 表达式读取(D46)、要能在事件里被安全拷贝——只有 dict 和 Pydantic 模型两条路能稳定满足这些。框架宁可在入口把不合规的类型拦下,也不愿意在持久化时才崩。🧠 今天你应该能回答
- Crew 编排不了哪三类需求?Flow 用什么模型解决?
- flow/ 为什么拆成 dsl / flow_definition / runtime 三层?中间的"契约"起什么作用?
- 公开
Flow类由哪两部分组合而成? - Flow 为什么继承 BaseModel?
ignored_types和FlowMeta各解决了什么麻烦? - 契约为什么懒构建 + 缓存?
- 一次 kickoff 是怎么"事件驱动"把方法图跑完的?
- Flow 与 Crew 是什么关系?谁编排谁?
✋ 10 分钟动手
P=lib/crewai/src/crewai/flow
sed -n '1,47p' $P/flow.py # 三层拆分 + 公开 Flow 组合
sed -n '428,476p' $P/runtime/__init__.py # 引擎 Flow 骨架 + 懒构建契约
sed -n '375,425p' $P/runtime/__init__.py # FlowMeta 元类
# 跑一个最小 Flow,亲眼看事件扩散
python -c "
from crewai.flow.flow import Flow, start, listen
class Demo(Flow):
@start()
def fetch(self): return 'raw'
@listen(fetch)
def clean(self, data): return data.upper()
print(Demo().kickoff()) # RAW
print(Demo.flow_definition().methods.keys()) # 看解析出的方法图
"
@start / @listen 把方法连成了图。明天钻进 flow/dsl/,逐行读这几个装饰器:它们怎么把一个普通函数包成 StartMethod、怎么把"我监听谁"记成 FlowMethodDefinition、or_/and_ 怎么表达"任一/全部"条件。