Day 41 / 共 60 天 · 阶段7 Flow 事件驱动

Flow 是什么:给 Agent 世界装一套"事件驱动的流程引擎"

前 40 天我们把 Agent / Task / Crew 拆到了底:Crew 会按 sequentialhierarchical 把一串任务跑完。但真实业务不是"一条直线"——有分支、有循环、有"等 A 和 B 都好了再做 C"。这类确定性流程正是 Flow 要解决的。今天先鸟瞰 crewai/flow/ 这个包:它被拆成 DSL(写法)/ Definition(契约)/ Runtime(引擎) 三层,公开的 Flow 类怎么拼出来,运行时骨架长啥样,以及一句 kickoff() 是怎么把方法按事件串起来跑的。

📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆 D41 Flow 总览 D42 装饰器 D43 定义契约 D44 路由跳转 D45 状态持久化 D46 表达式 D47 对话式 D48 选型 阶段8 LLM
💡 先用一个类比兜住今天 如果说 Crew 是"把一队人按顺序排好班让他们干活",那 Flow 就是一张流程图 + 一个自动派单的调度台。你在流程图上画好"哪个环节做完了,就通知下一个环节开工"(@start / @listen),调度台(Runtime)盯着每个环节完成的事件,一旦某个环节发出"我做完了"的信号,它就去查"谁在等这个信号",把等到的环节全叫起来跑。整个 Flow 的运转就是「方法完成 → 发事件 → 查谁在监听 → 触发下一批方法」的循环,而 Crew 只是流程图里可以被塞进某个环节的一块"积木"。
L01

痛点:Crew 为什么编排不了"带分支的流程"?

🤔 痛点Crew 很擅长"一串任务顺次跑完",但遇到这些需求就卡壳:① 分支——研究结果通过就发布,不通过就打回重写;② 循环——不满意就回到上一步再来一遍;③ 汇合——等"翻译"和"配图"都完成,才做"排版";④ 确定性——某几步我压根不想让 LLM 决策,就是写死的 Python 逻辑。这些是控制流问题,靠"任务列表 + Process"表达不出来。
💡 一句话本质 Flow 把"流程"建模成事件驱动的方法图:每个被 @start/@listen/@router 装饰的方法是图上的一个节点,节点之间靠"某方法完成"这个事件连边。kickoff() 启动后,引擎不停地做「执行方法 → 该方法完成的事件广播出去 → 找到监听这个事件的方法 → 执行它们」,直到没有方法可触发。控制流(分支/循环/汇合)由"谁监听谁"表达,业务逻辑就是普通 Python,LLM 只在你显式调用 Crew/Agent 时才介入。
大白话Crew 像"流水线":工位固定、顺序固定。Flow 像"发工单":谁干完了喊一声,系统自动把下一张工单派给正在等这个信号的人。你要做的只是在每个方法上贴一张"我等谁"的标签。
L02

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 这个老路径。
💡 设计取舍①:为什么要把"写法 / 契约 / 引擎"拆成三层? 朴素做法是一个大文件里既有装饰器又有执行逻辑。CrewAI 拆三层是因为:"怎么写"和"怎么跑"可以有多种组合。Python 装饰器只是产生 FlowDefinition 的一种方式——你也可以直接写 YAML 声明一个 Flow(Flow.from_declaration)。中间隔一层契约(Definition),写法层和引擎层就彻底解耦:换写法不影响引擎,换引擎不影响写法。代价是多了一次"翻译"(类 → Definition),换来的是可序列化、可可视化、可声明式三大能力。
数据结构:Flow 包的三层与数据流 写法层 dsl/ @start @listen @router / YAML 契约层 flow_definition FlowDefinition(可序列化的纯数据) 引擎层 runtime/ kickoff / 事件广播 / 状态 build_flow_definition 按契约执行 写法产生契约 → 引擎消费契约;中间隔一层,两端可各自替换
图注:装饰器(或 YAML)先生成 FlowDefinition,引擎再照着这份契约跑。契约是两端的"接口"。
L03

你 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 被约束为 dictBaseModel——这就是 Flow 的状态类型Flow[MyState] 让状态带上强类型。
start/listen/router/and_/or_装饰器和组合器也在这里重新导出,所以 from crewai.flow.flow import start 能用。
这种"公开类 = 引擎类 + Mixin"的写法是组合优于继承的变体:核心逻辑留在 RuntimeFlow,把可选的、实验性的对话能力做成一层可插拔的 Mixin,将来要去掉对话能力,改一行基类列表就行,不动引擎。
L04

运行时 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 缓存只算一次。"按需 + 缓存"是解析类元信息的标准优化。
L05

FlowMeta 元类:替你把"裸赋值"变成合法字段

🤔 痛点Flow 是 Pydantic 模型,可你在 Flow 类里经常写 some_config = 3 这种没有类型注解的裸赋值,或者放一些普通方法。纯 Pydantic 会因为"字段没注解"报错。谁来兜底?

答案是自定义元类 FlowMetaruntime/__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 能安全地被继承/扩展。
大白话元类就是"类的类",在你写的 Flow 类刚被 Python 造出来的那一刻插一脚,把那些"本来会让 Pydantic 生气"的写法悄悄修好。它让你写 Flow 时几乎感觉不到"我在写一个 Pydantic 模型"——这就是好框架的体贴。
L06

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)
控制流:kickoff 的事件驱动循环 kickoff() 执行 @start 方法 广播"完成"事件 查监听者 执行 @listen 方法 有人在等 它完成又发事件 → 继续扩散 没有任何方法能被触发时,整张图跑完,kickoff 返回最后一个方法的输出
图注:start 方法完成 → 发事件 → 找到监听者执行 → 监听者完成又发事件……事件像涟漪一样扩散,直到再没有可触发的方法。
📝 例子:最小 Flow 的一次 kickoff @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,流程却自己转完了。
L07

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", ...          # 人在环 / 输入
]
💡 关系一句话 Flow 编排,Crew 干活。Flow 负责"什么时候、按什么条件、以什么顺序"推进流程(控制流);Crew(以及单个 Agent)是流程里某一步被调用的执行单元。一个 Flow 方法里,你完全可以 result = MyCrew().kickoff(inputs=...),把 Crew 当普通函数用。反过来 Crew 里塞不进 Flow——因为 Crew 没有"事件监听"这套控制流能力。
维度CrewFlow
核心抽象Agent + Task 列表 + Process方法图 + 事件监听
控制流顺序 / 层级(写死两种)分支 / 循环 / 汇合(你画)
谁决策LLM 驱动为主你的 Python 代码驱动,LLM 按需介入
状态任务间 context 传递显式 self.state(可持久化)
定位被编排的"积木"顶层"编排者"

👶 小白:那我到底该用 Crew 还是 Flow?

👨‍🏫 老师:需求能用"一队 Agent 顺次/层级把任务做完"描述,就用 Crew,简单直接。一旦出现if/else 分支、要循环重试、要等多路汇合、要精确控制每一步,就上 Flow,并在需要"智能"的那一步里调用 Crew。大项目常见结构是:一个 Flow 编排全局,里面几个节点各自 kickoff 一个 Crew。D48 会专门讲选型。

L08

边界 + 今日小结

⚠️ 边界:状态类型只能是 dict 或 BaseModel,否则直接报错 Flow 的状态创建函数 _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_typesFlowMeta 各解决了什么麻烦?
  • 契约为什么懒构建 + 缓存?
  • 一次 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())   # 看解析出的方法图
"
明日预告 · Day 42:今天我们看到 @start / @listen 把方法连成了图。明天钻进 flow/dsl/,逐行读这几个装饰器:它们怎么把一个普通函数包成 StartMethod、怎么把"我监听谁"记成 FlowMethodDefinitionor_/and_ 怎么表达"任一/全部"条件。
← Day 40 embedding 存储 Day 42 · @start/@listen 装饰器 →