LlamaIndex Workflows 把 Agent 写成事件驱动的步骤网络:@step 装饰器声明步骤,事件在步骤间流动,Context 承载共享状态。本站教程 160 已经讲过它的基础编排——建工作流、串多步、接 RAG、管理状态。但官方文档里有一句话值得每个上生产的人停下来:「Workflows are ephemeral by default. Once run() returns, the state is gone.」默认状态下,一次 run 结束状态就没了,进程中途被杀,跑了一半的批量任务只能从头再来。
官方在 2026 年更新的 durable workflows 指南里给出了完整的持久化路径:在步骤边界把上下文快照成 JSON,崩溃后从快照恢复继续跑。这篇教程带你把这套机制落地:先跑通最小工作流,再实现手动快照与跨进程恢复,然后写一个自动检查点循环,接上并发扇出,最后亲手 kill 进程演练崩溃恢复。全部代码来自官方文档原文或其直接改写,每一步都可复现。
前置准备:Python 环境(建议 3.10 以上)与 pip。安装官方独立包:pip install llama-index-workflows。不需要任何模型 API Key——本篇的示例不调 LLM,专注在编排与持久化本身;接入真实 LLM 步骤时再配你自己的 Key。
LlamaIndex 生态有两套 import 路径:独立包 llama-index-workflows(from workflows import ...,官方当前主推,本篇采用)与旧的 llama_index.core.workflow(教程 160 所用,官方文档另有 WorkflowCheckpointer 等 API)。两套包的类名、事件名与内部结构随版本演进,代码跑不通时先核对官方文档与所装版本。
Step 1:跑通最小事件驱动工作流
用官方 PyPI 页的 quickstart 原文起步:两个步骤、一个自定义事件、一个带类型的共享状态。这个例子故意用纯 Python 逻辑不调模型,让「事件怎么流、状态怎么存」看得清清楚楚:
import asyncio
from pydantic import BaseModel, Field
from workflows import Context, Workflow, step
from workflows.events import Event, StartEvent, StopEvent
class MyEvent(Event):
msg: list[str]
class RunState(BaseModel):
num_runs: int = Field(default=0)
class MyWorkflow(Workflow):
@step
async def start(self, ctx: Context[RunState], ev: StartEvent) -> MyEvent:
async with ctx.store.edit_state() as state:
state.num_runs += 1
return MyEvent(msg=[ev.input_msg] * state.num_runs)
@step
async def process(self, ctx: Context[RunState], ev: MyEvent) -> StopEvent:
data_length = len("".join(ev.msg))
new_msg = f"Processed {len(ev.msg)} times, data length: {data_length}"
return StopEvent(result=new_msg)
async def main():
workflow = MyWorkflow()
ctx = Context(workflow)
result = await workflow.run(input_msg="Hello, world!", ctx=ctx)
print("Workflow result:", result)
asyncio.run(main())
跑一遍确认输出。注意 ctx.store.edit_state() 的写法——状态修改要在异步上下文管理器里进行,这是官方推荐的模式,快照机制依赖的正是这份状态。
Step 2:手动快照与跨进程恢复
官方 durable workflows 指南的核心模式只有三行骨架:run 拿到 handler,从 handler.ctx 导出快照,新进程里重建 Context 再续跑:
w = MyWorkflow()
handler = w.run()
result = await handler
# 持久化快照(示例存文件,生产建议落数据库)
import json
with open("my-run.json", "w") as f:
f.write(json.dumps(handler.ctx.to_dict()))
# ......新进程......
w2 = MyWorkflow()
with open("my-run.json") as f:
ctx = Context.from_dict(w2, json.loads(f.read()))
result = await w2.run(ctx=ctx) # 从恢复的状态继续
恢复语义官方文档写得很具体,值得逐条记住:恢复时会把「还在飞行中」的事件重新派发,重建部分的扇入缓冲;已经完成的步骤不会重跑(它们的输出已经在状态里);快照那一刻正在执行中的步骤会被回卷、从头再跑一遍。也就是说恢复是 at-least-once 的——同一个步骤的副作用可能发生不止一次。
at-least-once 恢复语义意味着所有步骤副作用必须幂等:写数据库用幂等键、发外部请求带请求 ID 去重、发邮件前查重。官方原文明确警告「step side effects need to be safe to repeat」——这不是建议,是恢复机制的前提条件。
Step 3:自动检查点循环:每个步骤完成就落盘
手动快照只保护「你记得存」的时刻。要扛住任意时刻的崩溃,官方给的检查点循环是监听内部事件——工作流每次步骤状态变化都会发内部事件,在步骤完成那一刻写快照:
from workflows.events import StepStateChanged, StepState
handler = w.run()
async for ev in handler.stream_events(expose_internal=True):
if isinstance(ev, StepStateChanged) and ev.step_state == StepState.NOT_RUNNING:
# 一个步骤刚跑完:此刻的状态是一致的检查点
save_snapshot("my-run.json", handler.ctx.to_dict())
result = await handler
官方同时提醒:繁忙的工作流不必在每个边界都写——每次快照是一次序列化加一次写盘,噪声大的工作流要节流,比如每 N 个检查点落一次盘或按时间间隔落盘。节流的代价要算清楚:崩溃后最多重做「上一个快照之后完成的工作」,间隔越长重做越多,这是一个用重做量换吞吐的旋钮。
给快照文件带版本号或时间戳再写入,恢复时挑最新的成功快照;写入用「写临时文件再原子改名」的方式,避免写到一半的快照被恢复逻辑读到。
Step 4:并发扇出与 Resource:让重活可重复、重资产不进快照
持久化最典型的受益场景是批量扇出:几百个文档逐个调模型,跑到一半进程死了,不应该从零再来。官方 durable 指南给的并发示例里有两个关键机制值得抄走——num_workers 控制并发度,Resource 注入的重型客户端(如 API client)不参与快照序列化,恢复时按需重建:
from typing import Annotated
from workflows.resource import Resource
class WorkItem(Event):
item_id: int
class WorkDone(Event):
item_id: int
result: str
def get_client():
return MyApiClient(...) # 重资产:连接、凭证都在这里
class EnrichBatch(Workflow):
@step
async def dispatch(self, ev: StartEvent) -> list[WorkItem]:
return [WorkItem(item_id=i) for i in ev.item_ids]
@step(num_workers=8)
async def enrich(self, ev: WorkItem,
client: Annotated[MyApiClient, Resource(get_client)]) -> WorkDone:
result = await client.enrich(ev.item_id) # 可重复的重活
return WorkDone(item_id=ev.item_id, result=result)
@step
async def collect(self, events: list[WorkDone]) -> StopEvent:
return StopEvent(result={ev.item_id: ev.result for ev in events})
把 Step 3 的检查点循环套在这个工作流上:每完成一个 enrich 就有快照。进程被杀后重启,从最新快照恢复,已完成的条目在状态里不会重跑,只有「死亡瞬间正在跑的那几个」会回卷重来——这正是 at-least-once 语义在批量场景下的具体表现。
num_workers 不是越大越好:并发度决定下游 API 的瞬时压力与费用速率。给慢而贵的外部调用留出限流余量,比崩溃后看着账单追悔划算;并发调参前先确认下游配额。
Step 5:崩溃演练:亲手 kill 一次再恢复
恢复代码没经历过真崩溃不算验证过。写一个长任务工作流(比如 dispatch 三十个 WorkItem,每个 enrich 里 sleep 一两秒模拟模型调用),启动后等它跑完三分之一,直接 kill 掉进程(Windows 用任务管理器结束进程树,macOS/Linux 用 kill 命令),确认快照文件是崩溃前最后一次检查点的内容。
重启进程,走 Step 2 的恢复路径:from_dict 重建 Context,run(ctx=ctx) 续跑,最终 collect 的结果应该与「不崩溃一口气跑完」完全一致。演练脚本的最小骨架:
# run_batch.py:有快照就恢复,没有就冷启动(对应官方 durable 指南原文结构)
async def run_batch(item_ids):
wf = EnrichBatch()
if os.path.exists(CKPT):
# 恢复:待办事件与扇入状态已在上下文里,
# 不要再发 StartEvent
ctx = Context.from_dict(wf, load_json(CKPT))
handler = wf.run(ctx=ctx)
else:
handler = wf.run(item_ids=item_ids)
await drive_with_checkpoints(handler) # Step 3 的检查点循环
多做几组对照:kill 得早、kill 得晚、kill 在并发高峰,观察「回卷重跑的条目集合」如何变化——你会直观看到节流间隔与重做量的关系。
恢复后的运行不要再发 StartEvent——官方示例明确注释「恢复后不要再发 StartEvent」,待办事件与扇入状态都已在恢复的上下文里,重复发起始事件等于把已完成的分派再做一遍。恢复分支与冷启动分支要在代码里显式分开。
Step 6:人工介入与生产配套
人工审批是持久化最自然的应用:跑到需要审批的步骤,发出「待审批」事件后让流程暂停,把此刻快照落库,等人工在后台系统点完通过,再从快照恢复并注入「审批通过」事件继续跑。审批可能等几小时甚至几天,靠内存挂着必然丢——快照进数据库,流程就成了可搁置、可迁移的实体。超时兜底同样在恢复层做:恢复逻辑检查快照时间戳,超时的审批走降级分支而不是无限等待。
生产配套三件事:可观测方面,Workflows 开箱即用地接了 OpenTelemetry,官方文档也点名了 Arize Phoenix——接上之后每个步骤的耗时与事件流都能看板化,崩溃点定位不再靠翻日志。评估方面,恢复演练要进 CI:用固定快照文件跑恢复断言,防止依赖升级悄悄破坏 from_dict 的兼容性。存储方面,快照里是完整的业务状态,按敏感数据对待:加密落库、访问控制、保留期一样不能少。
老包用户注意:llama_index.core.workflow.checkpointer 里有一个内建的 WorkflowCheckpointer,包装 run 方法后每个完成的步骤自动留检查点,还能用 filter_checkpoints 按步骤名筛选、run_from 从指定检查点续跑。新项目建议直接用本篇的独立包模式,存量代码迁移前先评估两套 API 的差异。
# 老包(llama_index.core.workflow)的内建检查点用法
from llama_index.core.workflow.checkpointer import WorkflowCheckpointer
ckptr = WorkflowCheckpointer(workflow=JokeFlow())
await ckptr.run(topic="chemistry")
# 按步骤名筛检查点,再从指定检查点续跑
after_gen = ckptr.filter_checkpoints(last_completed_step="generate_joke")[0]
await ckptr.run_from(checkpoint=after_gen)
预期效果自查:Step 1 输出「Processed 2 times, data length: 26」字样(状态在两次运行间累积);Step 2 的快照文件能被新进程读取并续跑成功;Step 5 的崩溃演练里,恢复后的最终结果与不崩溃对照组一致。三条全过,说明你的工作流已经具备「进程可死、任务不丢」的能力。
常见问题 FAQ
快照会很大吗?快照内容是事件队列加状态存储。Event 里挂大对象(整份文档、原始数据)会让快照膨胀且变慢——学官方示例用 Resource 装重型依赖,Event 只传 ID 与轻量结果,大文件放外部存储、快照里只存引用。
恢复时模型调用重跑一遍,多花的钱算谁的?算在你的节流策略头上。被回卷的步骤会真实重新执行,LLM 调用会再计费;把最贵的那一步放在「快照之后立刻完成」的位置,或对贵步骤单独做「执行前登记、执行后落结果」的幂等设计,能把重做成本压到接近于零。
和 Temporal 这类专用持久化引擎比呢?Workflows 的快照是轻量自托管方案,适合单机或自有基础设施里的中小规模任务;跨机器调度、长周期任务重试策略、版本化迁移这些重需求,Temporal 等引擎仍然更合适(本站教程 437 有专篇)。按任务规模选,不必二选一。