Skip to content

第 16 章 · Flow 进阶:状态管理与持久化

本章目标:掌握结构化与非结构化两种 state 的取舍,学会用 @persist 实现状态持久化与恢复(含 restore_from_state_id 分叉),用 or_ / and_ 编排多监听器的合流时机,用 @router 做条件路由分支。

16.1 两种 state 的再认识与操作细节

第 15 章已经见过字典式(unstructured)与 Pydantic 式(structured)state。这里补齐工程细节:

python
# 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),下次启动自动恢复:

python
# 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(方法级),实现细粒度控制:

python
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,新旧历史互不影响。
python
# 分叉示例
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_:多源监听的合流控制

当多个方法都可能触发同一个后续动作时,需要合流原语:

python
# 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() 让一个方法的输出变成"路标",把执行流引向不同的分支:

python
# 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 断点续跑的组合拳

把本章能力组合起来就是生产级的容错方案:

  1. 关键阶段方法挂 @persist(或类级持久化);
  2. 外部调用包 try/except,失败时通过 @router 把流程引入"降级分支"而不是直接崩;
  3. 进程重启后用 inputs={"id": ...} 续跑同一实例,或用 restore_from_state_id 从某个快照派生重放;
  4. 用第 14 章的事件监听(MethodExecutionStartedEvent / MethodExecutionFinishedEvent / FlowFailedEvent)记录每步落点,方便定位该从哪个快照恢复。
python
# 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,会发生什么?

🛠️ 动手实践

  1. 给第 15 章的写作流水线加 @persist,跑到一半 Ctrl-C 杀掉进程,再用 inputs={"id": ...} 续跑,验证 state 是否接续。
  2. 实现三路风控路由:金额 <1000 自动通过、<50000 走 AI 初审、其余人工;用 @router 串联两级路由并用 plot() 输出流程图。
  3. 构造一个"并行抓取两个数据源 + and_ 合流汇总"的 Flow,故意让其中一个数据源抛异常,观察流程卡住的现象,然后改造为"失败也返回标签"的降级版本。

工程化还差最后一环——脚手架与配置分离:第 17 章 · CLI 工程化与 YAML 配置项目