工作流
用 LlamaIndex Workflows 1.0 搭数据驱动的多步骤 Agent(事件驱动编排)
如果你的 Agent 任务是「数据密集、多步骤、要稳」——比如先检索文档、再总结、再校验、最后生成报告——直接写一串线性脚本很容易变成「面条代码」,一出错就全崩。LlamaIndex Workflows 1.0(2026 年 6 月发布)把编排从 RAG 库里独立成一套事件驱动原语:用 @step 把流程拆成可组合、可重试、可观测的步骤,步骤之间靠「事件」传递。本教程手把手带你从零搭一个数据驱动的多步骤 Agent。
📚 本教程适合:做 RAG / 文档处理 / 数据 pipeline 的 Python 开发者。你要会基础 Python、有一个 OpenAI(或兼容)Key。想做「对话型单 Agent」可看本中心 Dify / OpenAI Agents SDK 教程。
先搞懂:Workflow 和「普通脚本」差在哪?
用一句话理解:Workflow 是把「一个任务」拆成一串「步骤」,步骤靠事件通信,而不是靠函数层层嵌套调用。这带来三个好处:
| 维度 | 线性脚本 | LlamaIndex Workflow |
|---|---|---|
| 结构 | 函数嵌套,一环套一环 | @step 独立步骤,事件驱动 |
| 可组合 | 改一处牵全身 | 增删步骤互不影响 |
| 状态 | 全局变量乱传 | Context 显式管理共享状态 |
| 可观测 | 自己打日志 | 事件流天然可追踪 |
它和 LangGraph 的区别:Workflows 更轻、更偏「数据 pipeline + Agent 步骤」,没有显式状态图。如果你要复杂的分支/回滚/审计轨迹,LangGraph 更合适;要「把一堆步骤优雅串起来」,Workflows 更顺手。
Step 1:安装
1 装 llama-index 核心
pip install llama-index
# Workflows 是核心包的一部分,无需单独装
# 验证
python -c "from llama_index.core.workflow import Workflow; print('ok')"
# 设置模型 Key(以 OpenAI 为例)
export OPENAI_API_KEY=sk-xxxxxxxxxxxxxxxx
💡 LlamaIndex 把「检索层」和「编排层」分开了:RAG 用 LlamaIndex,多步骤编排用 Workflows。两者配合是 2026 年数据型 Agent 的主流组合。
Step 2:你的第一个 Workflow
2 用 @step + 事件跑起来
核心概念:StartEvent 是入口事件,StopEvent 是结束事件,@step 装饰的函数消费某类事件、产出下一类事件:
import asyncio
from llama_index.core.workflow import (
Workflow, StartEvent, StopEvent, step,
)
from llama_index.llms.openai import OpenAI
class MyFlow(Workflow):
@step
async def generate(self, ev: StartEvent) -> StopEvent:
llm = OpenAI(model="gpt-4o")
prompt = ev.get("topic", "AI Agent")
resp = await llm.acomplete(f"用一句话介绍:{prompt}")
return StopEvent(result=str(resp))
async def main():
flow = MyFlow(timeout=60, verbose=True)
result = await flow.run(topic="LlamaIndex Workflows")
print(result)
asyncio.run(main())
🔑 记住三件套:
StartEvent 进 → @step 处理 → StopEvent 出。ev.get("key") 取入口参数,verbose=True 能看到每一步事件流,调试神器。Step 3:多步骤串联
3 让步骤 A 的输出喂给步骤 B
自定义中间事件类型,把流程拆成「草稿 → 润色」两步:
from llama_index.core.workflow import (
Workflow, StartEvent, StopEvent, step,
)
from llama_index.llms.openai import OpenAI
from pydantic import BaseModel
class DraftEvent(StopEvent): # 中间事件(继承 Event 即可)
draft: str
class MyFlow(Workflow):
@step
async def draft(self, ev: StartEvent) -> DraftEvent:
llm = OpenAI(model="gpt-4o")
r = await llm.acomplete(f"写一句关于 {ev.get('topic')} 的标语")
return DraftEvent(draft=str(r))
@step
async def polish(self, ev: DraftEvent) -> StopEvent:
llm = OpenAI(model="gpt-4o")
r = await llm.acomplete(f"把这句标语改得更专业:{ev.draft}")
return StopEvent(result=str(r))
# run(topic="Workflows") → 先 draft 再 polish
步骤顺序由「事件类型」决定:某个 step 返回的事件类型,决定了下一个消费该事件的 step 会被触发。所以编排是「声明式」的——改事件连线就能改流程。
Step 4:接入 RAG 检索
4 让 Workflow 先查资料再回答
数据型 Agent 的灵魂就是「先检索再生成」。这里用内存向量索引做最小 RAG:
from llama_index.core import VectorStoreIndex, Document
from llama_index.core.workflow import Workflow, StartEvent, StopEvent, step
from llama_index.llms.openai import OpenAI
docs = [Document(text="RAG 是检索增强生成,先查资料再让模型回答。")]
index = VectorStoreIndex.from_documents(docs)
class RagFlow(Workflow):
@step
async def retrieve(self, ev: StartEvent):
retriever = index.as_retriever(similarity_top_k=1)
nodes = await retriever.aretrieve(ev.get("query"))
return QueryEvent(context="\n".join(n.get_content() for n in nodes))
@step
async def answer(self, ev: QueryEvent) -> StopEvent:
llm = OpenAI(model="gpt-4o")
r = await llm.acomplete(f"根据上下文回答:{ev.context}\n问题:{ev.query}")
return StopEvent(result=str(r))
# 记得定义 QueryEvent(BaseModel) 承载 context + query
💡 生产环境把
VectorStoreIndex.from_documents 换成持久化向量库(Qdrant / pgvector 等),Workflow 代码不变。这就是「LlamaIndex 管检索、Workflows 管编排」的分工。Step 5:共享状态与流式
5 用 Context 在步骤间存状态
当多个步骤需要共享中间结果(如累计 token、用户上下文),用 Context:
from llama_index.core.workflow import Workflow, StartEvent, StopEvent, step, Context
class Flow(Workflow):
@step
async def step1(self, ctx: Context, ev: StartEvent) -> StopEvent:
# 写共享状态
await ctx.set("user_id", ev.get("uid"))
# 读共享状态
uid = await ctx.get("user_id")
return StopEvent(result=f"用户 {uid} 的流程跑完了")
# 流式看进度(每步事件实时回调)
handler = flow.run(uid="u1")
async for ev in handler.stream_events():
print("事件:", type(ev).__name__)
Context vs 全局变量:Context 是每次运行隔离的,并发跑多个流程不会串数据,也比全局变量好测。多步骤共享状态一律走它。
Step 6:错误处理与可观测
6 超时、重试与追踪
1) 超时:构造时设 timeout,单步卡死会抛超时错误而非无限等
flow = MyFlow(timeout=120)
2) 重试:在 step 里用 try/except 包裹模型调用,
失败返回原事件类型即可让它「重新入队」重试
3) 可观测:Workflow 事件流可被 tracing 工具捕获
(如 LlamaIndex 自带 / OpenTelemetry),
哪一步慢、哪一步出错一目了然
4) 幂等与回放:因为步骤靠事件驱动,
可用事件日志重放整条流程做调试
💡 数据型 Agent 最怕「中间一步挂了全崩」。把模型调用单独包在 step 里 + 超时 + 重试,能把故障隔离在单步,而不是拖垮整条 pipeline。
Step 7:落地为服务
7 包成 API 对外提供
# 用 FastAPI 把 Workflow 暴露成 HTTP 接口
from fastapi import FastAPI
app = FastAPI()
flow = MyFlow()
@app.post("/run")
async def run(topic: str):
return {"result": await flow.run(topic=topic)}
# 启动:uvicorn main:app --reload
# 调用:curl -X POST localhost:8000/run?topic=Agent
🎉 到此你已跑通 Workflows 核心:建步骤 → 多步串联 → RAG 检索 → 共享状态/流式 → 容错观测 → 服务化。它最适合「数据密集、要稳、要可追溯」的多步骤 Agent。
常见问题速查
| 现象 | 大概率原因 & 解决 |
|---|---|
| 流程不往下走 | 没有 step 返回对应事件类型 → 检查返回的事件类是否匹配 |
| 超时报错 | 某步模型调用太慢 → 调大 timeout 或加重试 |
| Context 取不到值 | 先 set 再 get,且在同一 run 内;别跨 run 共享 |
| RAG 检索为空 | 文档没建索引 / top_k 太小 / query 与文档语义不匹配 |