FlowDefinition:Flow 的"设计图纸"长什么样
Day 42 每个装饰器都往方法上挂了一份 FlowMethodDefinition。今天把镜头拉远:这些方法定义汇总成一份完整的 FlowDefinition——Flow 的可序列化契约。我们读它的字段结构(schema/name/state/config/methods)、三个把关的校验器、build_flow_definition 怎么扫描 Flow 类把图纸还原出来,以及 flow_config 全局配置和 from_declaration——后者让你连一行 Python 都不写,光靠 YAML 就声明出一个能跑的 Flow。
FlowDefinition 就是一栋楼的建筑图纸:图纸上标了每个房间(method)、房间之间的门(listen 条件)、承重规范(config)、地基类型(state)。图纸不是楼——但施工队(runtime 引擎)照着图纸就能盖出楼。图纸还能存进档案馆(序列化成 YAML)、能打印出来给人看(可视化)、能被审图员挑错(校验器)。今天读的就是这张图纸的格式和"从代码逆向出图纸"的过程。痛点:装饰器已经能跑了,为什么还要一层"契约"?
FlowDefinition?多这一层不是绕吗?FlowDefinition 就是这个纯数据契约:它把"图长什么样"从"用什么语言写的图"里剥离出来。装饰器只是生成契约的一种方式;YAML 是另一种。引擎只认契约,不认你怎么写的。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
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流程级的执行配置、持久化、对话开关。都是可选、带默认。三个校验器:图纸出厂前的"审图"
契约是 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),拼错的表达式在构建期就报错,不等到运行时。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_definition(flow/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 的硬校验。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_config(flow/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() # 进程级单例,部署时在启动处覆盖
| FlowConfigDefinition | flow_config(单例) | |
|---|---|---|
| 作用域 | 单个 Flow 契约 | 整个进程 |
| 内容 | stream/max_method_calls 等执行参数 | HITL provider / 输入 provider |
| 谁设 | Flow 类的字段 / YAML | 部署代码在启动时设一次 |
| 典型用途 | "这个 Flow 要不要流式" | "本次部署人在环走 WebSocket 而非控制台" |
max_method_calls=100 是重要的安全阀:Flow 支持循环(D44),万一你写了个死循环,引擎跑满 100 次方法调用会强制停下,避免无限烧钱。凡是支持循环的引擎,都要有这类硬上限。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}) # 把契约注入实例
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 表达式。边界 + 今日小结
@model_validator 会 raise(硬错误,Flow 直接建不出来),而 L05 的 log_flow_definition_issues(flow/flow_definition.py:960)只打日志不 raise(软警告,Flow 还能跑)。二者的分界是:"这个问题会不会让流程必然出错?" 语法非法、监听自己 → 必错,硬拦;监听了一个当前没定义但可能来自路由的事件 → 未必错,只警告。坑在于:软警告很容易被忽略。如果你的 listener 死活不触发,先去翻日志里有没有 "flow definition issue" —— 很可能是事件名拼错了,但它只是警告没报错。🧠 今天你应该能回答
- 装饰器已经能跑了,为什么还要抽一层 FlowDefinition?它换来哪四种能力?
FlowMethodDefinition的do为什么不止是"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))
"
router/emit 字段却没细讲。明天专攻 @router 与条件跳转:router 的返回值怎么变成"下一个事件名"、引擎的 _execute_listeners 怎么处理路由结果、以及"分支/循环"在事件模型里到底是怎么实现的。