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

FlowDefinition:Flow 的"设计图纸"长什么样

Day 42 每个装饰器都往方法上挂了一份 FlowMethodDefinition。今天把镜头拉远:这些方法定义汇总成一份完整的 FlowDefinition——Flow 的可序列化契约。我们读它的字段结构(schema/name/state/config/methods)、三个把关的校验器、build_flow_definition 怎么扫描 Flow 类把图纸还原出来,以及 flow_config 全局配置和 from_declaration——后者让你连一行 Python 都不写,光靠 YAML 就声明出一个能跑的 Flow

📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆 D41 Flow 总览 D42 装饰器 D43 定义契约 D44 路由跳转 D45 状态持久化 D46 表达式 D47 对话式 D48 选型 阶段8 LLM
💡 先用一个类比兜住今天 FlowDefinition 就是一栋楼的建筑图纸:图纸上标了每个房间(method)、房间之间的门(listen 条件)、承重规范(config)、地基类型(state)。图纸不是楼——但施工队(runtime 引擎)照着图纸就能盖出楼。图纸还能存进档案馆(序列化成 YAML)、能打印出来给人看(可视化)、能被审图员挑错(校验器)。今天读的就是这张图纸的格式和"从代码逆向出图纸"的过程。
L01

痛点:装饰器已经能跑了,为什么还要一层"契约"?

🤔 痛点Day 42 的装饰器已经把图信息挂到方法上了,引擎似乎可以直接读方法来跑。为什么 CrewAI 还要专门抽出一份 FlowDefinition?多这一层不是绕吗?
💡 一句话本质 因为"能跑"只是最低要求,Flow 还想要可序列化、可校验、可可视化、可声明式四种能力,而这四种能力都需要一个脱离 Python 类的纯数据表示FlowDefinition 就是这个纯数据契约:它把"图长什么样"从"用什么语言写的图"里剥离出来。装饰器只是生成契约的一种方式;YAML 是另一种。引擎只认契约,不认你怎么写的。
大白话把"设计"和"实现"分开:图纸(Definition)描述要盖什么,施工队(Runtime)负责盖。图纸能存档、能审、能画给人看——这些是"一堆装饰过的方法"给不了的。
L02

FlowMethodDefinition:一个方法在图纸上的全部信息

先看单个方法的定义(flow/flow_definition.py:643):

# flow/flow_definition.py:643
class FlowMethodDefinition(BaseModel):
    """Static definition of one Flow method and its execution roles."""
    description: str | None = None                    # 人读的说明
    do: FlowActionDefinition = Field(...)             # 这方法执行什么动作(code/crew/agent/...)
    start: bool | FlowDefinitionCondition | None = None  # 是否起点 / 起点条件
    listen: FlowDefinitionCondition | None = None     # 监听条件(等谁)
    router: bool = False                              # 是否路由方法(输出当事件名)
    emit: list[str] | None = None                     # 声明可能发出的路由事件
    human_feedback: FlowHumanFeedbackDefinition | None = None  # 人在环步骤
    persist: FlowPersistenceDefinition | None = None  # 方法级持久化覆盖
do★核心:方法"做什么"。不止是 Python 代码——FlowActionDefinition 可以是 code / crew / agent / expression / tool / script 多种(这也是为什么 Flow 能被声明式描述)。
start / listenDay 42 那两个装饰器填的就是这两个字段。互斥语义:起点填 start,监听者填 listen。
router + emitD44 主角:router=True 表示"我的返回值当下一个事件名",emit 声明可能的取值(供可视化和静态校验)。
human_feedback / persist可选增强:人在环、方法级持久化。它们由 @human_feedback/@persist 装饰器填。

它还有一个巧妙的规范化校验器(flow/flow_definition.py:689):

# flow/flow_definition.py:689
@model_validator(mode="after")
def _canonicalize_human_feedback_routing(self) -> FlowMethodDefinition:
    # 一个 human_feedback 声明了 emit 结果的方法,就按 router 对待
    if self.human_feedback is not None and self.human_feedback.emit:
        self.router = True
        self.emit = None
    return self
这体现了契约的一个原则:无论用户怎么写,都归一到"标准形态"。带 emit 的人在环方法本质就是"人来路由"(人点了"通过/打回")——那就把它标成 router,让引擎用同一套路由逻辑处理,而不是给它再写一套特例。规范化 = 把多种写法收敛成一种内部表示,引擎只需处理一种。
L03

FlowDefinition:整张图纸的顶层结构

把所有方法定义装进来的顶层契约(flow/flow_definition.py:710):

# flow/flow_definition.py:710
class FlowDefinition(BaseModel):
    schema_: Literal["crewai.flow/v1"] = Field(default="crewai.flow/v1", alias="schema")
    name: str = Field(...)                            # 流程名,日志/事件/trace 里用
    description: str | None = None
    state: FlowStateDefinition | None = None          # 状态契约(dict / pydantic / json_schema)
    config: FlowConfigDefinition = Field(default_factory=FlowConfigDefinition)
    persist: FlowPersistenceDefinition | None = None  # 流程级持久化
    conversational: FlowConversationalDefinition | None = None  # 对话式配置(D47)
    methods: dict[str, FlowMethodDefinition] = Field(default_factory=dict)  # ★方法名 → 定义
schema_ 带版本★契约自带 schema: crewai.flow/v1 版本号。将来格式升级到 v2,可据此兼容老 YAML。alias="schema" 因为 schema 是 Pydantic 保留名。
methods: dict[名字→定义]★图纸的主体:一张"方法名 → 方法定义"的表。谁监听谁,全在各 method 的 listen 字段里;顶层不存边,边是分布式记在节点上的。
state状态的类型契约。dict 型、pydantic 型、json_schema 型三种——决定 kickoff 时怎么造初始状态(D45)。
config / persist / conversational流程级的执行配置、持久化、对话开关。都是可选、带默认。
数据结构:FlowDefinition 的字段树 FlowDefinition schema v1 name state config methods: dict "fetch" → MethodDef do/start/listen/router 边(谁监听谁)不在顶层,而是分散记在每个 MethodDef 的 listen 里
图注:FlowDefinition 是一棵字段树,核心是 methods 表;每个 method 自带 listen 条件,边是"分布式"存储的。
L04

三个校验器:图纸出厂前的"审图"

契约是 Pydantic 模型,构建后自动跑三个校验器(flow/flow_definition.py:767):

# flow/flow_definition.py:767
@model_validator(mode="after")
def _validate_method_names(self) -> FlowDefinition:
    for method_name in self.methods:
        _validate_step_name(method_name, field="Flow method names")   # ① 名字合法
    return self

@model_validator(mode="after")
def _validate_trigger_namespace(self) -> FlowDefinition:
    for method_name, method in self.methods.items():
        if _condition_references(method.listen, method_name):         # ② 不能监听自己
            raise ValueError(f"methods.{method_name}.listen must not reference itself")
    return self

@model_validator(mode="after")
def _validate_cel_expressions(self) -> FlowDefinition:
    for method_name, method in self.methods.items():
        _validate_action_cel(method.do, path=f"methods.{method_name}.do",   # ③ 表达式合法
                             allowed_roots=_BASE_CEL_ROOTS)
    return self
① 名字合法方法名要能当"事件名/标识符"用,非法字符会被拦。
② 不能监听自己D42 提过的死循环防护——@listen("a") def a 在这被拒。
③ CEL 表达式合法★如果 do 是 expression 动作(D46),这里静态检查表达式只引用了允许的根(state/outputs),拼错的表达式在构建期就报错,不等到运行时。
💡 设计取舍①:为什么把校验放在"构建契约时"而不是"运行时"? 朴素做法是跑到哪校验到哪。但 Flow 常常是长流程 + 调 LLM/Crew(慢且花钱)。如果一个拼错的表达式要等流程跑了 10 分钟、烧了几美元 token 才在第 8 步炸出来,体验极差。CrewAI 选择fail-fast:在契约构建/加载的瞬间就把能静态发现的错全查一遍——名字、自引用、表达式根。代价是构建稍慢一点点,换来的是"错误尽早、尽廉价地暴露"。对慢而贵的流程,前置校验的性价比极高。
L05

build_flow_definition:从 Python 类逆向出图纸

装饰器把信息挂在方法上,_build_flow_definition_from_class 负责把它们收集成完整契约(flow/dsl/_utils.py:435):

# flow/dsl/_utils.py:435
def _build_flow_definition_from_class(flow_class, namespace=None) -> FlowDefinition:
    methods: dict[str, FlowMethodDefinition] = {}
    flow_methods = _iter_flow_methods(flow_class)          # ① 扫出所有被装饰的方法
    ...
    for method_name, method in flow_methods.items():
        methods[method_name] = _build_method_definition(   # ② 每个方法 → 一份 MethodDef
            method, f"methods.{method_name}")

    definition = FlowDefinition(                           # ③ 组装顶层契约
        name=getattr(flow_class, "__name__", "Flow"),
        description=(flow_class.__doc__ or "").strip() or None,
        state=_build_state_definition(flow_class),         # 从 Flow[T] 的 T 推状态契约
        config=_build_config_definition(flow_class),
        persist=_build_persistence_definition(flow_class),
        conversational=_build_conversational_definition(flow_class),
        methods=methods,
    )
    log_flow_definition_issues(definition)                 # ④ 记录(非致命的)诊断
    return definition

_build_method_definitionflow/dsl/_utils.py:361)把装饰器挂的碎片和 do 动作合成一份:

# flow/dsl/_utils.py:361
def _build_method_definition(method, path) -> FlowMethodDefinition:
    fragment = _get_flow_method_definition(method)         # 取装饰器挂的元数据
    if fragment is None:
        method_definition = FlowMethodDefinition(do=_method_action(method))
    else:
        method_definition = fragment.model_copy(deep=True, update={"do": _method_action(method)})
    human_feedback = _build_human_feedback_definition(method, ...)  # 合入人在环
    if human_feedback is not None:
        method_definition.human_feedback = human_feedback
    method_definition.persist = _build_persistence_definition(method)  # 合入持久化
    return method_definition
_iter_flow_methods遍历类的属性,用 is_flow_method 挑出被装饰过的(是 FlowMethod 实例的)。还会处理继承、对话式方法、和 Pydantic 字段名冲突的特殊情况。
_build_state_definition★从 Flow[MyState] 的类型参数 MyState 推出状态契约:是 dict 就记 dict 型,是 BaseModel 就记 pydantic 型 + import ref。
fragment + do装饰器只挂了 start/listen/router 等,没挂 do(因为要用类名限定)。build 时补上 do=_method_action,合成完整 MethodDef。
log_flow_definition_issues★非致命问题(如某方法监听了不存在的事件)在这记日志,但不 raise——区别于 L04 的硬校验。
控制流:build_flow_definition 的流水线 Flow 类 (装饰过的方法) _iter_flow_methods 扫出被装饰方法 每方法→MethodDef 碎片+do 合成 3 个校验器 fail-fast 缓存 ClassVar 懒执行:首次 flow_definition() 才跑,结果按类缓存,之后直接取
图注:从 Flow 类到可用契约的五步流水线——扫描、逐方法合成、审图、缓存;整条流水线懒执行、只算一次。
大白话这一步是"从盖好的样板房逆向画出图纸":挨个房间(方法)看它头上贴的便利贴(装饰器元数据),抄进图纸;再补上"这房间实际是干嘛的"(do 动作),最后拼成整张图。
L06

FlowConfigDefinition 与 flow_config:两种"配置"别搞混

流程级执行配置(flow/flow_definition.py:191):

# flow/flow_definition.py:191
class FlowConfigDefinition(BaseModel):
    tracing: bool | None = None
    stream: bool = False                    # 是否流式发事件
    memory: dict[str, Any] | None = None
    input_provider: str | None = None
    suppress_flow_events: bool = False      # 关掉事件(内部 Flow 用)
    max_method_calls: int = 100             # ★一次 kickoff 最多执行多少次方法(防跑飞)
    defer_trace_finalization: bool = False
    checkpoint: bool | dict[str, Any] | None = None

另有一个全局单例 flow_configflow/flow_config.py:17),是完全不同的东西:

# flow/flow_config.py:17
class FlowConfig:
    """Global configuration for Flow execution."""
    def __init__(self) -> None:
        self._hitl_provider: HumanFeedbackProvider | None = None   # 全局人在环 provider
        self._input_provider: InputProvider | None = None          # 全局 ask() 输入源
    ...
flow_config = FlowConfig()   # 进程级单例,部署时在启动处覆盖
FlowConfigDefinitionflow_config(单例)
作用域单个 Flow 契约整个进程
内容stream/max_method_calls 等执行参数HITL provider / 输入 provider
谁设Flow 类的字段 / YAML部署代码在启动时设一次
典型用途"这个 Flow 要不要流式""本次部署人在环走 WebSocket 而非控制台"
max_method_calls=100 是重要的安全阀:Flow 支持循环(D44),万一你写了个死循环,引擎跑满 100 次方法调用会强制停下,避免无限烧钱。凡是支持循环的引擎,都要有这类硬上限。
L07

from_declaration:一行 Python 都不写,纯 YAML 造 Flow

契约既然是纯数据,就能反过来"从数据造 Flow"(flow/flow_definition.py:809):

# flow/flow_definition.py:809
@classmethod
def from_declaration(cls, *, contents=None, path=None) -> FlowDefinition:
    if isinstance(contents, cls):
        return contents
    source_path = None
    if contents is None:
        if path is None:
            raise ValueError("Provide contents or path")
        source_path = Path(path)
        contents = source_path.expanduser().read_text(encoding="utf-8")  # 读 YAML 文件
    if isinstance(contents, dict):
        return cls._load_mapping(contents)
    ...
    loaded = yaml.safe_load(contents)                # 解析 YAML → dict
    if not isinstance(loaded, dict):
        raise ValueError("Flow declaration must contain a mapping")
    return cls._load_mapping(loaded, source_path=source_path)  # dict → 校验 → FlowDefinition

运行时的 Flow.from_declaration 再把契约变成可跑的实例(runtime/__init__.py:478):

# runtime/__init__.py:478
@classmethod
def from_declaration(cls, *, contents=None, path=None, **kwargs) -> Flow[Any]:
    definition = FlowDefinition.from_declaration(contents=contents, path=path)
    return cls.model_validate(
        {**definition.config.model_dump(), **kwargs},
        context={"flow_definition": definition})    # 把契约注入实例
📝 例子:一份最小 YAML Flow
schema: crewai.flow/v1
name: ResearchFlow
state:
  type: dict
  default: {topic: "AI agents"}
methods:
  seed:
    start: true
    do: {call: expression, expr: state.topic}
  greet:
    listen: seed
    do: {call: expression, expr: "'研究: ' + state.topic"}
Flow.from_declaration(path="research.yaml").kickoff() —— 没写一个 Python 方法,Flow 照跑。do 里的 expression 就是 D46 的 CEL 表达式。
💡 设计取舍②:为什么要支持"声明式 YAML Flow"? Python 写法灵活但需要程序员。声明式 YAML 让非开发者(产品、运营)也能配流程,也让流程能存进数据库、由平台可视化编辑、热更新而不发版。代价是 YAML 里只能用受限的动作(expression/crew/agent 等,不能写任意 Python),表达力弱于代码。这正是 D41 那层"契约"的终极价值兑现:因为契约是纯数据,写法就能有 Python 和 YAML 两套,各取所需。
L08

边界 + 今日小结

⚠️ 边界:契约构建失败 vs 诊断日志,两条不同的线 注意 L04 的三个 @model_validatorraise(硬错误,Flow 直接建不出来),而 L05 的 log_flow_definition_issuesflow/flow_definition.py:960)只打日志不 raise(软警告,Flow 还能跑)。二者的分界是:"这个问题会不会让流程必然出错?" 语法非法、监听自己 → 必错,硬拦;监听了一个当前没定义但可能来自路由的事件 → 未必错,只警告。坑在于:软警告很容易被忽略。如果你的 listener 死活不触发,先去翻日志里有没有 "flow definition issue" —— 很可能是事件名拼错了,但它只是警告没报错。

🧠 今天你应该能回答

  • 装饰器已经能跑了,为什么还要抽一层 FlowDefinition?它换来哪四种能力?
  • FlowMethodDefinitiondo 为什么不止是"Python 代码"?
  • Flow 的"边"(谁监听谁)存在顶层还是各方法里?
  • 三个校验器分别拦什么?为什么校验放在构建期?
  • build_flow_definition 大致分哪几步?
  • FlowConfigDefinition 和全局 flow_config 有什么区别?max_method_calls 防什么?
  • from_declaration 让 Flow 多了什么能力?它成立的前提是什么?

✋ 10 分钟动手

P=lib/crewai/src/crewai/flow
sed -n '643,708p'  $P/flow_definition.py       # FlowMethodDefinition + 规范化
sed -n '710,795p'  $P/flow_definition.py       # FlowDefinition + 三校验器
sed -n '361,383p'  $P/dsl/_utils.py            # 每方法 → MethodDef
sed -n '435,469p'  $P/dsl/_utils.py            # 组装顶层契约
# 把一个 Flow 类的契约打印成 YAML 看看
python -c "
from crewai.flow.flow import Flow, start, listen
class Demo(Flow):
    @start()
    def a(self): return 1
    @listen(a)
    def b(self, x): return x+1
import json; print(json.dumps(Demo.flow_definition().to_dict(), indent=2, ensure_ascii=False))
"
明日预告 · Day 44:今天见到了 router/emit 字段却没细讲。明天专攻 @router 与条件跳转:router 的返回值怎么变成"下一个事件名"、引擎的 _execute_listeners 怎么处理路由结果、以及"分支/循环"在事件模型里到底是怎么实现的。
← Day 42 装饰器 Day 44 · router 与条件跳转 →