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:拷贝 + 补 iddict(init_state) 先拷一份——直接改原 dict 会污染"类上的默认状态",两次 kickoff 就串了。
BaseModel 有 iddump 成 dict,缺 id 补 id,再重建实例。保证每次 kickoff 状态干净独立。
BaseModel 无 id → 动态混入★用户的状态模型忘了继承 FlowState?框架当场造一个 StateWithId 子类把 id 补上——多重继承的实用技巧,用户少写一行也能持久化。
兜底 TypeError(D41 提过) 既不是 dict 也不是 BaseModel → 直接报错。持久化/序列化只认这两种。
数据结构:state 的两种形态,都带 id dict 型 { "id": "a1b2...", "topic": "AI", "count": 3 } 灵活、无类型检查 Pydantic 型 (FlowState) id: str ← 来自 FlowState topic: str = "AI" count: int = 0 有类型/校验/补全,推荐
图注:两种状态形态都保证有 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 见过这个标志)
控制流:存档与读档 方法A完成 INSERT 一行(id最新) flow_states 表 (WHERE uuid=X) …,方法A,{...} ← 最新 …,方法B(旧),{...} kickoff(inputs={"id":X}) load_state 取最新 跳过已完成 → 接着跑
图注:每个方法完成 INSERT 一行;恢复时按 uuid 取最新一行、进入恢复模式跳过已完成方法。
L08

边界 + 今日小结

⚠️ 边界:state 必须能序列化成 JSON,塞进不可序列化对象会存不下 持久化时 state 要经 _to_state_dictmodel_dump()/转 dict)再 json.dumps 存进 state_jsonflow/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。
← Day 44 路由跳转 Day 46 · expressions 表达式 →