1. LlamaIndex工作流:从入门到实战
作为一名长期从事AI应用开发的工程师,我深知构建复杂大语言模型(LLM)应用时的痛点。当系统需要调用多次模型、查询多个数据源、处理分支逻辑时,传统的顺序代码很快就会变得难以维护。这就是为什么当我发现LlamaIndex Workflows这个轻量级事件驱动框架时,立刻被它的设计理念所吸引。
LlamaIndex Workflows于2025年6月发布1.0正式版,作为一个独立的Python/TypeScript包,它提供了一种全新的方式来构建AI应用。其核心理念是将复杂流程拆解为多个独立的"步骤"(Step),步骤之间通过"事件"(Event)通信,由框架负责调度执行。这种设计带来了显著的优点:代码模块化、易于测试、支持并行执行,还内置了可视化调试工具。
在接下来的内容中,我将通过大量实际代码示例,带你全面掌握这个强大的编排框架。无论你是想构建RAG机器人、多智能体协作系统,还是需要人工审批的内容生成流水线,LlamaIndex Workflows都能成为你的得力助手。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 核心组件解析
2.1 事件(Event)基础
事件是LlamaIndex Workflows中最基本的通信单元。所有事件都继承自Event类,它本质上是一个Pydantic模型,可以携带结构化数据。让我们看一个简单的例子:
python复制from llama_index.core.workflow import Event
from typing import List, Optional
# 定义携带简单消息的事件
class MessageEvent(Event):
content: str
# 定义更复杂的事件
class AnalysisEvent(Event):
topic: str
keywords: List[str]
confidence: float
除了自定义事件,框架还提供了两个特殊的内置事件:
- StartEvent:工作流的入口事件,调用
workflow.run()时传入的参数会自动封装成它 - StopEvent:工作流的出口事件,当某个步骤返回它时,工作流会立即终止
python复制from llama_index.core.workflow import StartEvent, StopEvent
# StartEvent可以携带任意字段
# StopEvent需要传入result参数
2.2 工作流类定义
定义好事件后,我们需要创建一个继承自Workflow的类,并在其中定义各个步骤。每个步骤都用@step装饰器标记:
python复制from llama_index.core.workflow import Workflow, step
from llama_index.llms.openai import OpenAI
class MyWorkflow(Workflow):
# 可以初始化共享资源,如LLM实例
llm = OpenAI(model="gpt-4o-mini")
@step
async def first_step(self, ev: StartEvent) -> MessageEvent:
# 处理逻辑...
return MessageEvent(content="处理完成")
框架会根据方法的参数类型注解和返回值类型注解,自动判断该步骤接收什么事件、产出什么事件。
2.3 上下文(Context)管理
当工作流变得复杂时,单纯靠事件传递数据会非常繁琐。Context对象就像一个全局的"数据黑板",任何步骤都可以随时读写数据:
python复制from llama_index.core.workflow import Context
class StatefulWorkflow(Workflow):
@step
async def first_step(self, ctx: Context, ev: StartEvent) -> None:
# 写入数据到上下文
await ctx.set("user_query", ev.query)
await ctx.set("start_time", datetime.now())
@step
async def second_step(self, ctx: Context, ev: SomeEvent) -> StopEvent:
# 从上下文读取数据
query = await ctx.get("user_query")
start_time = await ctx.get("start_time")
print(f"处理查询: {query}, 耗时: {datetime.now() - start_time}")
return StopEvent(result="完成")
需要注意的是,ctx.get()和ctx.set()都是异步方法,需要使用await。
3. 高级功能与技巧
3.1 多事件等待与并发
有时候,一个步骤需要等待多个事件全部到达后才能执行。Context提供了collect_events()方法来实现这一点:
python复制class UserProfileEvent(Event):
profile: str
class RecommendationEvent(Event):
items: List[str]
class AggregationWorkflow(Workflow):
@step
async def fetch_profile(self, ctx: Context, ev: StartEvent) -> UserProfileEvent:
await asyncio.sleep(1)
return UserProfileEvent(profile="科技爱好者")
@step
async def fetch_recommendations(self, ctx: Context, ev: StartEvent) -> RecommendationEvent:
await asyncio.sleep(1)
return RecommendationEvent(items=["GPU", "机械键盘"])
@step
async def aggregate(self, ctx: Context, ev: StartEvent) -> StopEvent:
events = await ctx.collect_events(
ev, [UserProfileEvent, RecommendationEvent]
)
if events is None:
return None # 等待更多事件
profile_event, rec_event = events
result = f"为用户 {profile_event.profile} 推荐 {rec_event.items}"
return StopEvent(result=result)
3.2 流式事件处理
对于LLM生成这种长耗时操作,实时反馈进度能极大提升用户体验:
python复制class ProgressEvent(Event):
msg: str
class StreamingWorkflow(Workflow):
llm = OpenAI(model="gpt-4o-mini")
@step
async def generate(self, ctx: Context, ev: StartEvent) -> StopEvent:
ctx.write_event_to_stream(ProgressEvent(msg="开始生成..."))
full_response = ""
async for chunk in self.llm.astream_complete(ev.prompt):
full_response += chunk.delta
ctx.write_event_to_stream(ProgressEvent(msg=chunk.delta))
ctx.write_event_to_stream(ProgressEvent(msg="生成完成!"))
return StopEvent(result=full_response)
3.3 检查点与恢复
对于长时间运行的工作流,检查点(Checkpoint)机制允许保存完整状态并在以后恢复:
python复制class CheckpointWorkflow(Workflow):
@step
async def critical_step(self, ctx: Context, ev: StartEvent) -> NextEvent:
result = await self.do_something()
await self.save_checkpoint(ctx, "after_critical_step")
return NextEvent(data=result)
# 使用检查点
workflow = CheckpointWorkflow()
handler = workflow.run()
try:
result = await handler
except Exception as e:
last_checkpoint = workflow.get_last_checkpoint()
new_handler = workflow.run_from(checkpoint=last_checkpoint)
result = await new_handler
4. 实战案例:智能客服系统
让我们通过一个完整的智能客服案例来展示LlamaIndex Workflows的实际应用:
python复制from enum import Enum
from datetime import datetime
class QuestionType(Enum):
AFTER_SALES = "售后"
PRE_SALES = "售前"
COMPLAINT = "投诉"
class CustomerServiceWorkflow(Workflow):
@step
async def classify(self, ctx: Context, ev: StartEvent) -> ClassifyEvent:
logs = await ctx.get("logs", default=[])
logs.append(f"[{datetime.now()}] 收到问题: {ev.question}")
await ctx.set("logs", logs)
# 模拟分类逻辑
question = ev.question.lower()
if "退货" in question or "维修" in question:
qtype = QuestionType.AFTER_SALES
elif "多少钱" in question or "价格" in question:
qtype = QuestionType.PRE_SALES
else:
qtype = QuestionType.COMPLAINT
return ClassifyEvent(qtype=qtype, question=ev.question)
@step
async def handle_after_sales(self, ctx: Context, ev: ClassifyEvent) -> AfterSalesEvent:
if ev.qtype != QuestionType.AFTER_SALES:
return None
logs = await ctx.get("logs")
logs.append(f"[{datetime.now()}] 进入售后流程")
await ctx.set("logs", logs)
return AfterSalesEvent(question=ev.question)
@step
async def generate_response(self, ctx: Context, ev: AfterSalesEvent) -> StopEvent:
response = f"【售后】关于「{ev.question}」,请提供订单号,我们将为您安排退货/维修。"
logs = await ctx.get("logs")
logs.append(f"[{datetime.now()}] 生成回复: {response}")
await ctx.set("logs", logs)
print("\n=== 处理日志 ===")
for log in logs:
print(log)
print("===============\n")
return StopEvent(result=response)
5. 调试与可视化
LlamaIndex提供了强大的可视化工具来帮助理解和调试工作流:
python复制from llama_index.utils.workflow import (
draw_all_possible_flows,
draw_most_recent_execution,
)
# 绘制静态流程图
draw_all_possible_flows(MyWorkflow, filename="workflow_structure.html")
# 绘制最近一次执行的动态轨迹
workflow = MyWorkflow()
await workflow.run(input="测试")
draw_most_recent_execution(workflow, filename="recent_execution.html")
动态执行图特别有用,可以清楚地看到哪些分支被实际执行了,哪一步耗时最长。
6. 部署方案
工作流可以轻松部署为Web服务。以下是TypeScript版本的示例:
typescript复制import { Hono } from "hono";
import { createHonoHandler } from "@llamaindex/workflow-core/interrupter/hono";
const app = new Hono();
app.post("/api/run", createHonoHandler(
myWorkflow,
async (ctx) => startEvent(await ctx.req.json()),
stopEvent
));
serve(app);
Python版本也可以使用FastAPI等框架进行类似封装。
7. 经验总结与最佳实践
在实际项目中使用LlamaIndex Workflows几个月后,我总结出以下几点经验:
- 模块化设计:将复杂逻辑拆分为小而专的步骤,每个步骤只做一件事
- 合理使用Context:避免过度依赖全局状态,只将真正需要共享的数据放入Context
- 错误处理:为关键步骤添加检查点,确保长时间运行的工作流可以恢复
- 监控与日志:利用流式事件和可视化工具监控工作流执行情况
- 测试策略:由于步骤是独立的,可以轻松为每个步骤编写单元测试
一个常见陷阱是在Context中存储过多数据,这会导致状态管理混乱。我的建议是:如果数据只在相邻步骤间传递,优先使用事件;只有真正全局需要的数据才放入Context。
另一个实用技巧是使用stepwise=True模式调试复杂工作流:
python复制async def debug_workflow():
workflow = MyComplexWorkflow()
handler = workflow.run(stepwise=True)
# 逐步执行
events = await handler.run_step()
print(f"第一步产出的事件: {events}")
events = await handler.run_step()
print(f"第二步产出的事件: {events}")
final_result = await handler
print(f"最终结果: {final_result}")
这种模式下,工作流不会自动运行到底,而是每执行一个步骤就暂停,方便调试。
