第 16 章 · Flow 进阶:状态管理与持久化
本章目标:掌握结构化与非结构化两种 state 的取舍,学会用
@persist实现状态持久化与恢复(含restore_from_state_id分叉),用or_/and_编排多监听器的合流时机,用@router做条件路由分支。
16.1 两种 state 的再认识与操作细节
第 15 章已经见过字典式(unstructured)与 Pydantic 式(structured)state。这里补齐工程细节:
# state_ops.py —— 状态操作的完整对照
from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, start
class OrderState(BaseModel):
order_id: str = ""
amount: float = 0.0
risk_level: str = ""
class UnstructuredFlow(Flow):
@start()
def load(self):
# 字典式:键随意增删,自动带 id 键(UUID)
print(f"State ID: {self.state['id']}")
self.state["order_id"] = "A-1024"
self.state["amount"] = 199.0
class StructuredFlow(Flow[OrderState]):
@start()
def load(self):
# Pydantic 式:属性访问,字段拼错直接报错
self.state.order_id = "A-1024"
self.state.amount = 199.0
print(f"State ID: {self.state.id}") # id 同样自动维护选择口径(官方建议):结构简单或高度动态、原型期 → 字典;需要一致的结构、类型安全、IDE 校验 → Pydantic。团队协作项目默认选 Pydantic——它把"字段名写错"这类 bug 从运行时深处提前到了编码时。
16.2 @persist:让 Flow 断电可续
长流水线跑一半进程崩溃/重启,从头再来既费钱又费时。@persist 装饰器把 state 自动存入本地 SQLite(默认后端 SQLiteFlowPersistence),下次启动自动恢复:
# persist_flow.py —— 类级持久化
from pydantic import BaseModel
from crewai.flow.flow import Flow, start, listen
from crewai.flow.persistence import persist
class CounterState(BaseModel):
counter: int = 0
@persist # 类级:所有方法的状态都会持久化
class CounterFlow(Flow[CounterState]):
@start()
def step(self):
self.state.counter += 1
print(f"[id={self.state.id}] counter={self.state.counter}")
# 第一次运行:counter 0 -> 1,快照写入 SQLite
CounterFlow().kickoff()也可以只在关键方法上挂 @persist(方法级),实现细粒度控制:
class AnotherFlow(Flow[dict]):
@persist # 只持久化这一个方法的状态
@start
def begin(self):
if "runs" not in self.state:
self.state["runs"] = 0
self.state["runs"] += 1
print("已持久化的运行次数:", self.state["runs"])恢复语义分两种(官方文档明确区分):
- 续跑(resume):
kickoff(inputs={"id": <uuid>})—— 加载该 UUID 的最新快照,继续在同一个flow_uuid下追加历史; - 分叉(fork):
kickoff(restore_from_state_id=<uuid>)—— 用旧快照初始化新运行的 state,但分配新的state.id,新旧历史互不影响。
# 分叉示例
flow_1 = CounterFlow()
flow_1.kickoff() # counter -> 1
flow_2 = CounterFlow()
flow_2.kickoff(restore_from_state_id=flow_1.state.id)
# flow_2 从 counter=1 起步,随后 step() 使其变为 2;
# flow_2.state.id 是新 ID,flow_1 的历史不受影响。持久化注意事项
- 结构化与字典式 state 都支持;
id字段缺失会自动补上; - 若
restore_from_state_id找不到对应快照,kickoff 会静默回退为全新运行——排查"为什么没恢复"时要先确认 UUID 正确; - 它与
from_checkpoint不能同时使用(会抛ValueError),二选一。
16.3 or_ 与 and_:多源监听的合流控制
当多个方法都可能触发同一个后续动作时,需要合流原语:
# or_and_flow.py —— or_ 触发任意一个,and_ 等齐所有
from crewai.flow.flow import Flow, listen, start, or_, and_
class AlertFlow(Flow):
@start()
def check_cpu(self):
return "CPU 正常"
@listen(check_cpu)
def check_disk(self):
self.state["disk"] = "磁盘告警"
return "磁盘告警"
@listen(or_(check_cpu, check_disk)) # 任一完成即触发(可能触发多次)
def logger(self, result):
print(f"[OR] 日志: {result}")
@listen(and_(check_cpu, check_disk)) # 全部完成后才触发一次
def summary(self):
print(f"[AND] 汇总: {self.state['disk']} / CPU 已检查")
return "巡检完成"
flow = AlertFlow()
print(flow.kickoff())运行输出能看到 [OR] 日志 打印了两次(每次源方法完成都触发一次),而 [AND] 汇总 只在两个来源都完成后执行一次。经验法则:
- 日志/审计类旁路动作用
or_(每个事件都要记); - 汇总/决策类动作用
and_(必须等齐所有输入); and_的监听方法不接收单个返回值参数(多个来源无法映射成一个入参),数据请走 state。
16.4 @router:条件路由分支
@router() 让一个方法的输出变成"路标",把执行流引向不同的分支:
# router_flow.py —— 按风险等级路由审批流
import random
from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, router, start
class RiskState(BaseModel):
amount: float = 0.0
success_flag: bool = False
class ApprovalFlow(Flow[RiskState]):
@start()
def submit_order(self):
self.state.amount = random.choice([99.0, 99000.0])
# 大额订单标记高风险(真实场景这里调用风控服务)
self.state.success_flag = self.state.amount < 10000
@router(submit_order) # 返回值就是路由标签
def route_by_risk(self):
if self.state.success_flag:
return "auto_approve"
return "manual_review"
@listen("auto_approve")
def approve(self):
print(f"订单 {self.state.amount} 自动通过")
return "approved"
@listen("manual_review")
def review(self):
print(f"订单 {self.state.amount} 进入人工审核队列")
return "needs_human"
if __name__ == "__main__":
ApprovalFlow().plot("approval_plot") # plot 能看到分支图
print(ApprovalFlow().kickoff())要点:@router(被监听方法) 的函数体是普通 Python 判断逻辑,返回的字符串决定走哪条路;下游用 @listen("标签") 接住。配合 plot() 生成的 HTML 图,复杂分支一目了然。多级路由可以串联:分支方法本身还可以再被 @router 监听,形成决策树。
16.5 断点续跑的组合拳
把本章能力组合起来就是生产级的容错方案:
- 关键阶段方法挂
@persist(或类级持久化); - 外部调用包 try/except,失败时通过
@router把流程引入"降级分支"而不是直接崩; - 进程重启后用
inputs={"id": ...}续跑同一实例,或用restore_from_state_id从某个快照派生重放; - 用第 14 章的事件监听(
MethodExecutionStartedEvent/MethodExecutionFinishedEvent/FlowFailedEvent)记录每步落点,方便定位该从哪个快照恢复。
# resume_demo.py —— 失败降级 + 可恢复的最小骨架
import os
from pydantic import BaseModel
from crewai.flow.flow import Flow, listen, router, start
from crewai.flow.persistence import persist
class PipelineState(BaseModel):
step_done: int = 0
@persist
class RobustFlow(Flow[PipelineState]):
@start()
def stage1(self):
self.state.step_done = 1 # 完成即持久化
return "s1_ok"
@listen(stage1)
def stage2(self):
try:
raise TimeoutError("模拟外部服务超时")
except TimeoutError:
self.state.step_done = 2
return "s2_failed" # 返回失败标签而非抛出
@router(stage2)
def route(self):
return "retry_stage2" # 真实场景可查 state 决定重试或跳过
@listen("retry_stage2")
def on_retry(self):
print(f"从快照恢复后可重试,当前进度: 第 {self.state.step_done} 阶段已完成")
RobustFlow().kickoff()
# 进程崩溃后:
# RobustFlow().kickoff(inputs={"id": "<上次打印的 state.id>"})16.6 本章小结
- 字典 state 胜在灵活,Pydantic state 胜在类型安全与校验,团队项目默认后者;两者都自动携带唯一
id; @persist支持类级与方法级,默认 SQLite 后端;inputs={"id":...}是同实例续跑,restore_from_state_id是分叉出新 ID 的重放;or_任一完成即触发(可能多次),and_全部完成才触发一次且数据要走 state;@router把方法返回字符串变成路由标签,@listen("标签")承接分支,plot()可视化决策树;- 生产容错 = 持久化 + 降级分支 + 快照恢复 + 事件留痕的组合拳。
🧪 随堂测验
点击你认为正确的选项。答错时会展示正确答案与原因解析。
1. 关于 and_(a, b) 监听方法的触发行为,正确的是?
2. kickoff(restore_from_state_id=<uuid>) 的语义是?
3. @router 装饰的方法靠什么决定走哪个分支?
4. 若 restore_from_state_id 传入了一个不存在的 uuid,会发生什么?
🛠️ 动手实践
- 给第 15 章的写作流水线加
@persist,跑到一半 Ctrl-C 杀掉进程,再用inputs={"id": ...}续跑,验证 state 是否接续。 - 实现三路风控路由:金额 <1000 自动通过、<50000 走 AI 初审、其余人工;用
@router串联两级路由并用plot()输出流程图。 - 构造一个"并行抓取两个数据源 +
and_合流汇总"的 Flow,故意让其中一个数据源抛异常,观察流程卡住的现象,然后改造为"失败也返回标签"的降级版本。
工程化还差最后一环——脚手架与配置分离:第 17 章 · CLI 工程化与 YAML 配置项目。