教程中心工作流
工作流

用 LlamaIndex Workflows 1.0 搭数据驱动的多步骤 Agent(事件驱动编排)

2026.07.23· 7 个步骤 · 19 分钟阅读· 📚 LlamaIndex

如果你的 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 与文档语义不匹配