Day 45 / 共 60 天 · 阶段7 Flow 事件驱动
state 与持久化:让流程"能记住自己跑到哪了"
Day 44 的 router 一直在读 self.state。今天把状态这条线彻底讲透:状态对象是什么(FlowState 自带 id)、self.state 怎么初始化(dict 还是 Pydantic)、@persist 怎么让"每个方法完成后自动落盘"、落盘走的 FlowPersistence 抽象接口、内置的 SQLiteFlowPersistence 怎么存怎么取,以及 kickoff(inputs={"id": ...}) 怎么"接着上次那次跑"。这是 Flow 能做长流程、可恢复、人在环的地基。
📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆→
D41 Flow 总览→
D42 装饰器→
D43 定义契约→
D44 路由跳转→
D45 状态持久化→
D46 表达式→
D47 对话式→
D48 选型→
阶段8 LLM
💡 先用一个类比兜住今天
state 就是流程的存档进度,@persist 就是游戏的自动存档。每通关一个关卡(方法完成),游戏自动把进度写进存档栏(SQLite 一行)。哪天断电(进程崩),你重开游戏、读那个存档栏的 id,就能从上次通关的地方接着打,而不用从头再来。存档栏的编号就是 state.id——认准它,就能找回任意一次流程的进度。L01
痛点:流程跑到第 8 步崩了,前 7 步白跑?
🤔 痛点Flow 常常是长流程:调好几个 Crew、每步都花钱花时间。如果第 8 步网络抖了一下崩了,难道前 7 步(可能跑了 5 分钟、烧了几美元)全作废、从头再来?还有"人在环"——流程停下等人审批,人第二天才来,这期间进程早关了,怎么恢复?
💡 一句话本质
答案是把状态持久化:每个方法完成,就把当前
self.state 序列化存进数据库,键是 state.id。恢复时用同一个 id load_state 回来,跳过已完成的方法,从断点接着跑。state 是"当下的全部进度",持久化是"把进度写下来",id 是"进度的门牌号"。三者合起来,长流程才敢崩、才敢停下等人。大白话没有持久化的流程像"不能存档的游戏",一崩全没。
@persist 给它加了自动存档,崩了、关了都能读档继续。L02
FlowState 与 self.state:状态的最小共识是"有个 id"
所有结构化状态的基类,只强制一件事:有唯一 id(runtime/__init__.py:360):
# runtime/__init__.py:360
class FlowState(BaseModel):
"""Base model for all flow states, ensuring each state has a unique ID."""
id: str = Field(
default_factory=lambda: str(uuid4()), # 不给就自动生成一个 UUID
description="Unique identifier for the flow state")
# runtime/__init__.py:1659 你写 self.state 时读到的
@property
def state(self) -> T:
return cast(T, self._state)
FlowState 只有一个 id★持久化的最小共识:不管你的状态多复杂,必须有个 id 当持久化键。你的 Pydantic 状态继承 FlowState 就自动有了。default_factory=uuid4没显式给 id,就自动生成一个 UUID。所以你 MyFlow().kickoff() 时其实已经有一个隐形的 state.id 了。state 是 propertyself.state 读的是私有属性 _state(D41 讲过,不是模型字段)。T 是你在 Flow[MyState] 里指定的类型。📝 例子:两种状态写法
无类型:
有类型:
class F(Flow): ... → state 是 dict,自动塞一个 {"id": "..."}。有类型:
class S(BaseModel): topic: str = "",class F(Flow[S]) → state 是 S 实例,框架自动让它带上 id(继承 FlowState)。推荐后者:字段有类型、IDE 有补全、CEL 表达式好读。L03
_create_initial_state:dict 和 BaseModel 分别怎么造
状态初始化要兼顾多种写法,核心分支(runtime/__init__.py:1531):
# runtime/__init__.py:1588(节选)
if isinstance(init_state, dict):
new_state = dict(init_state) # 拷贝一份,避免改到类属性
if "id" not in new_state:
new_state["id"] = str(uuid4()) # dict 也要补 id
return cast(T, new_state)
if isinstance(init_state, BaseModel):
model = init_state
if hasattr(model, "id"):
state_dict = model.model_dump()
if not state_dict.get("id"):
state_dict["id"] = str(uuid4())
return cast(T, type(model)(**state_dict))
# 模型没 id → 动态混入 FlowState 补上 id
class StateWithId(FlowState, type(model)):
pass
state_dict = model.model_dump(); state_dict["id"] = str(uuid4())
return cast(T, StateWithId(**state_dict))
dict:拷贝 + 补 id★dict(init_state) 先拷一份——直接改原 dict 会污染"类上的默认状态",两次 kickoff 就串了。BaseModel 有 iddump 成 dict,缺 id 补 id,再重建实例。保证每次 kickoff 状态干净独立。BaseModel 无 id → 动态混入★用户的状态模型忘了继承 FlowState?框架当场造一个 StateWithId 子类把 id 补上——多重继承的实用技巧,用户少写一行也能持久化。兜底 TypeError(D41 提过) 既不是 dict 也不是 BaseModel → 直接报错。持久化/序列化只认这两种。图注:两种状态形态都保证有 id;持久化时统一
model_dump() 或直接取 dict 存成 JSON。L04
@persist:又是一个"只贴标签"的装饰器
@persist 本身不干活,只往类/方法上盖个持久化配置的章(flow/persistence/decorators.py:147):
# flow/persistence/decorators.py:183
def decorator(target: type | Callable[..., T]) -> type | Callable[..., T]:
actual_persistence = (
persistence if persistence is not None else default_flow_persistence() # 默认 SQLite
)
_stamp_persistence_metadata(target, actual_persistence, verbose) # 盖章
return target
# :56 盖的章长这样
def _stamp_persistence_metadata(target, persistence, verbose) -> None:
target.__flow_persistence_config__ = SimpleNamespace(
persistence=persistence, verbose=verbose)
docstring 明说"pure metadata stamper"★源码原话:"The decorator is a pure metadata stamper"。它不改行为,只记配置,真正存盘由引擎驱动。和 D42 的装饰器一个思路。类级 / 方法级盖在 class 上 → 所有方法都存;盖在单个方法上 → 只存那个方法完成时。方法级配置优先于类级。默认 default_flow_persistence()不指定后端就用默认(通常 SQLite)。你也可以 @persist(MyPostgresPersistence())。SimpleNamespace轻量地把 persistence + verbose 打包成一个对象挂上去,引擎日后读它。和 D43 呼应:
_build_persistence_definition 会把这个章读进 FlowMethodDefinition.persist / FlowDefinition.persist,成为契约的一部分——所以持久化配置也能被序列化、被 YAML 声明。L05
引擎:每个方法完成后调 _persist_method_completion
D41 的 _execute_method 里,方法一完成就落盘(runtime/__init__.py:2610):
# runtime/__init__.py:2608
self._completed_methods.add(method_name)
await asyncio.to_thread(self._persist_method_completion, method_name) # ★存盘(放线程池,不阻塞事件循环)
# runtime/__init__.py:2668
def _persist_method_completion(self, method_name) -> None:
method_definition = self._definition.methods[str(method_name)]
persist_definition = (
method_definition.persist # 方法级优先
if method_definition.persist is not None
else self._definition.persist) # 否则用流程级
if persist_definition is None or not persist_definition.enabled:
return # 没配持久化 → 直接返回,什么都不存
from crewai.flow.persistence.decorators import PersistenceDecorator
backend = (self.persistence
if self._instance_persistence and self.persistence is not None
else self._persist_backend_for(persist_definition))
PersistenceDecorator.persist_state(self, method_name, backend, verbose=persist_definition.verbose)
to_thread(...)★存盘是同步 IO(SQLite 写文件)。丢进线程池跑,避免堵住 asyncio 事件循环,其它 listener 还能并行。方法级 or 流程级先看这方法自己有没有 persist 配置,没有才用流程级。粒度可控。没配就 return不加 @persist 的 Flow,这里直接返回——零开销,不存的流程不为持久化付一分钱代价。backend 选择实例上显式传的 persistence= 优先于契约里派生的后端——运行时能覆盖声明。💡 存档时机为什么是"每个方法完成后"?因为方法是 Flow 的最小可恢复单元。方法执行到一半崩了没法"半存",但一个方法完整跑完、结果进了 state,这个点就是干净的存档点。恢复时以"已完成方法集合"为准,跳过它们、从下一个继续。选"方法边界"当存档点,是可恢复性和开销的平衡。
L06
FlowPersistence:可换后端的抽象接口
持久化后端是抽象基类,SQLite / Postgres / 内存都实现它(flow/persistence/base.py:18):
# flow/persistence/base.py:18
class FlowPersistence(BaseModel, ABC):
persistence_type: str = Field(default="base")
def __init_subclass__(cls, **kwargs) -> None: # 子类自动注册
super().__init_subclass__(**kwargs)
if not getattr(cls, "__abstractmethods__", set()):
_persistence_registry[cls.__name__] = cls
@abstractmethod
def init_db(self) -> None: ... # 建表/连库
@abstractmethod
def save_state(self, flow_uuid, method_name, state_data) -> None: ... # 存
@abstractmethod
def load_state(self, flow_uuid) -> dict[str, Any] | None: ... # 取最新
# 非抽象:async 人在环用(默认实现 = 直接存/不支持)
def save_pending_feedback(self, flow_uuid, context, state_data) -> None:
self.save_state(flow_uuid, context.method_name, state_data)
ABC + abstractmethod三个必实现方法:建库、存、取。这就是"持久化后端"的最小契约。__init_subclass__ 自动注册★你一定义一个新后端子类,它就自动进 _persistence_registry——之后能按类名从注册表反查(配合 YAML 声明式指定后端)。save/load 以 flow_uuid 为键就是 state.id。存取都围着它转。pending_feedback 有默认实现★人在环相关方法非抽象,给了默认(退化成普通存)。意思是"后端可选支持"——普通后端不实现也能用,想支持异步人在环再覆盖。💡 设计取舍①:为什么把持久化做成抽象接口 + 注册表,而不是写死 SQLite?
写死 SQLite 最省事,但生产环境千差万别:开发用 SQLite 文件、测试用内存假实现、生产可能要 Postgres 或公司自研的存储服务。抽象成
FlowPersistence 接口,配合 set_flow_persistence_factory 全局工厂(D43 的 default_flow_persistence),就能在部署启动处一次性换掉后端,业务代码里的 @persist 一个字不改。代价是多一层抽象、多一个注册机制。"面向接口 + 依赖注入"换来的是环境无关性——框架该有的觉悟。L07
SQLiteFlowPersistence:一行 INSERT + 取最新
建表(flow/persistence/sqlite.py:71)——注意它是追加式历史表:
# flow/persistence/sqlite.py:78
conn.execute("PRAGMA journal_mode=WAL") # 并发友好
conn.execute("""
CREATE TABLE IF NOT EXISTS flow_states (
id INTEGER PRIMARY KEY AUTOINCREMENT,
flow_uuid TEXT NOT NULL, -- = state.id
method_name TEXT NOT NULL, -- 哪个方法完成时存的
timestamp DATETIME NOT NULL,
state_json TEXT NOT NULL -- 整个 state 序列化成 JSON
)""")
conn.execute("CREATE INDEX IF NOT EXISTS idx_flow_states_uuid ON flow_states(flow_uuid)")
取状态时取同 uuid 最新的一行(flow/persistence/sqlite.py:178):
-- flow/persistence/sqlite.py:189
SELECT state_json FROM flow_states
WHERE flow_uuid = ?
ORDER BY id DESC -- 自增 id 最大 = 最近一次
LIMIT 1
追加式历史表★每次 save 是 INSERT 新行,不是 UPDATE 覆盖。同一个 flow_uuid 会有多行——每个方法完成一行。这就天然记录了完整历史。ORDER BY id DESC LIMIT 1load 时取自增 id 最大的那行 = 最近一次状态。历史保留,但恢复用最新。WAL 模式 + store_lockSQLite 的 WAL 日志模式 + 进程锁,让多个方法(线程池里)并发写不打架。state_jsonstate 序列化成 JSON 文本存。这就是为什么 state 必须是 dict/Pydantic——它们能稳定转 JSON。恢复流程:kickoff(inputs={"id": 已有uuid}) 时,引擎 load_state 回来、进入"恢复模式",跳过已完成方法(runtime/__init__.py:2066):
# runtime/__init__.py:2066
is_restoring = (inputs and "id" in inputs and self.persistence is not None) \
or self._restored_from_checkpoint
if not is_restoring:
self._completed_methods.clear() # 全新执行:清空进度
self._method_outputs.clear(); ...
else:
if self._completed_methods:
self._is_execution_resuming = True # ★恢复模式:已完成的方法会被跳过(D44 见过这个标志)
图注:每个方法完成 INSERT 一行;恢复时按 uuid 取最新一行、进入恢复模式跳过已完成方法。
L08
边界 + 今日小结
⚠️ 边界:state 必须能序列化成 JSON,塞进不可序列化对象会存不下
持久化时
state 要经 _to_state_dict(model_dump()/转 dict)再 json.dumps 存进 state_json(flow/persistence/sqlite.py:170)。如果你往 state 里塞了个不能序列化的东西——比如一个数据库连接、一个 Crew 实例、一个 lambda——存盘时就会炸,或悄悄丢失。persist_state 里对存盘异常包了 RuntimeError("State persistence failed: ...")(flow/persistence/decorators.py:132)。正确姿势:state 里只放"数据"(字符串、数字、列表、嵌套 Pydantic),不放"活对象"。要用 Crew,在方法里临时 new 一个、用完把结果(纯数据)写进 state,别把 Crew 本身挂到 state 上。另一个坑:state 缺 id 会直接 ValueError("Flow state must have an 'id' field")——但正常继承 FlowState 就不会遇到。🧠 今天你应该能回答
- 没有持久化的 Flow 有什么问题?id 在持久化里扮演什么角色?
FlowState为什么强制一个 id?dict 状态的 id 从哪来?- 用户状态模型忘了继承 FlowState,框架怎么补 id?
@persist改变行为吗?真正的存盘由谁驱动、在什么时机?- 存盘为什么放
asyncio.to_thread?不加 @persist 的流程有开销吗? FlowPersistence做成抽象接口 + 注册表,换来什么?- SQLite 表为什么是"追加式"?load 怎么取到最新状态?恢复时靠什么跳过已完成方法?
✋ 10 分钟动手
P=lib/crewai/src/crewai/flow
sed -n '1531,1611p' $P/runtime/__init__.py # _create_initial_state
sed -n '2668,2690p' $P/runtime/__init__.py # _persist_method_completion
sed -n '147,192p' $P/persistence/decorators.py # @persist 贴标签
sed -n '71,203p' $P/persistence/sqlite.py # 建表 / save / load
# 跑一个带持久化的 Flow,然后用同一个 id 恢复
python -c "
from pydantic import BaseModel
from crewai.flow.flow import Flow, start, listen
from crewai.flow.persistence import persist
class S(BaseModel):
id: str=''
step: int=0
@persist()
class F(Flow[S]):
@start()
def a(self): self.state.step=1; return 'a'
@listen(a)
def b(self, _): self.state.step=2; return 'b'
f=F(); f.kickoff(); print('id=',f.state.id,'step=',f.state.step)
"
明日预告 · Day 46:D43 的 YAML Flow 里
do: {call: expression, expr: state.topic} 用了 CEL 表达式。明天读 expressions.py:${...} 模板怎么解析、CEL 怎么安全地从 state/outputs 里取值、为什么用 CEL 而不是直接 eval Python。