一、核心概念速览LlamaIndex Workflows 是事件驱动、可断点续跑、支持人机协同的结构化工作流引擎用于编排复杂的 LLM 任务RAG、Agent、多步骤数据处理等。核心能力事件驱动通过事件流转控制流程全局上下文跨步骤共享状态 / 数据检查点崩溃后可恢复执行人机协同流程暂停等待人工输入多事件等待同时等待多个触发条件逐步执行调试时单步运行可部署生产环境稳定运行二、环境准备pip install llama-index-core llama-index-llms-openai python-dotenv三、完整工作流代码实现覆盖所有需求点1. 定义工作流事件事件是流程的流转信号可携带数据。from llama_index.core.workflow import ( Workflow, Event, StartEvent, StopEvent, step, Context, ) from typing import Optional, List, Dict, Any import asyncio # 自定义事件 class DataLoadedEvent(Event): 数据加载完成事件 data: List[Dict[str, Any]] class LLMProcessedEvent(Event): LLM 处理完成事件 result: str class HumanReviewEvent(Event): 人工审核事件人机协同 approve: bool feedback: Optional[str] None class ErrorEvent(Event): 错误事件 message: str2. 定义全局上下文 / 状态全局上下文用于跨步骤共享数据支持检查点持久化。# 全局状态类检查点会自动序列化 class WorkflowState: def __init__(self): self.raw_data: List[Dict] [] self.llm_result: str self.human_approved: bool False self.human_feedback: str self.process_log: List[str] []3. 定义工作流主类核心集成入口点、退出、多事件等待、手动触发、人机协同、检查点、逐步执行。class AIAssistantWorkflow(Workflow): AI 助手工作流全功能演示 def __init__(self, timeout: int 300): super().__init__(timeouttimeout) # 初始化全局状态 self.global_state WorkflowState() # 工作流入口点 step async def start_step(self, ctx: Context, ev: StartEvent) - DataLoadedEvent | ErrorEvent: 入口初始化流程 加载数据 try: self.global_state.process_log.append(进入工作流入口) # 模拟加载数据 self.global_state.raw_data [{content: 用户需求总结 LlamaIndex 工作流}] return DataLoadedEvent(dataself.global_state.raw_data) except Exception as e: return ErrorEvent(messagef初始化失败{str(e)}) # 核心业务步骤 step async def llm_process_step(self, ctx: Context, ev: DataLoadedEvent) - HumanReviewEvent: LLM 处理步骤 self.global_state.process_log.append(开始 LLM 处理) # 模拟 LLM 调用 self.global_state.llm_result LlamaIndex Workflows 是事件驱动的 AI 工作流引擎支持检查点、人机协同... # 触发人工审核 return HumanReviewEvent(approveFalse, feedback等待人工审核) # 人机协同等待手动触发 step async def human_review_step(self, ctx: Context, ev: HumanReviewEvent) - LLMProcessedEvent | ErrorEvent: 人工审核流程暂停等待手动触发事件 self.global_state.process_log.append(流程暂停等待人工审核) # 保存检查点 await ctx.save_checkpoint() # 等待人工手动触发审核结果多事件等待 human_ev await ctx.wait_for_event(HumanReviewEvent) # 更新全局状态 self.global_state.human_approved human_ev.approve self.global_state.human_feedback human_ev.feedback if human_ev.approve: self.global_state.process_log.append(人工审核通过) return LLMProcessedEvent(resultself.global_state.llm_result) else: return ErrorEvent(messagef审核不通过{human_ev.feedback}) # 多事件等待可选扩展 async def wait_multiple_events(self, ctx: Context): 同时等待多个事件示例 ev1, ev2 await ctx.wait_for_any(LLMProcessedEvent, ErrorEvent) return ev1 or ev2 # 工作流退出 step async def end_step(self, ctx: Context, ev: LLMProcessedEvent | ErrorEvent) - StopEvent: 退出返回最终结果 if isinstance(ev, ErrorEvent): final_result {status: failed, message: ev.message} else: final_result { status: success, data: self.global_state.llm_result, log: self.global_state.process_log } # 清空临时状态 self.global_state.process_log.append(工作流结束) return StopEvent(resultfinal_result)4. 工作流调用手动触发、逐步执行、检查点async def main(): # 1. 初始化工作流 workflow AIAssistantWorkflow(timeout600) # 2. 启动工作流非阻塞 run_task asyncio.create_task(workflow.run()) # 3. 等待流程执行到人工审核步骤 await asyncio.sleep(1) print( 当前工作流状态 ) print(等待人工审核...) # 4. 手动触发事件人机协同核心 await workflow.send_event( HumanReviewEvent(approveTrue, feedback内容合格通过审核) ) # 5. 获取最终结果 result await run_task print(\n 工作流最终结果 ) print(result) # 逐步执行调试用 # workflow AIAssistantWorkflow() # async for step_result in workflow.run_stream(): # print(f步骤执行结果{step_result}) # 检查点恢复崩溃后续跑 # 加载检查点上下文后直接调用await workflow.resume_from_checkpoint(ctx) if __name__ __main__: asyncio.run(main())四、工作流流程图可视化五、核心知识点详解1. 工作流入口 / 退出入口step接收StartEvent即为入口点退出返回StopEvent(result...)即为退出result是最终输出2. 全局上下文 / 状态协同定义状态类存储共享数据步骤内直接修改 / 读取实现跨步骤数据传递检查点会自动序列化状态支持断点续跑3. 多事件等待# 等待任意一个事件 ev await ctx.wait_for_any(EventA, EventB) # 等待所有事件 evs await ctx.wait_for_all(EventA, EventB)4. 手动事件触发# 外部手动向工作流发送事件 await workflow.send_event(自定义事件(参数值))5. 人机协同流程中调用ctx.wait_for_event()实现暂停外部手动触发事件后流程自动恢复执行适合审核、确认、修正等人工干预场景6. 逐步执行# 流式逐步执行每一步返回结果 async for step in workflow.run_stream(): print(step)7. 检查点工作流# 保存检查点 await ctx.save_checkpoint() # 从检查点恢复执行 await workflow.resume_from_checkpoint(ctx)六、生产环境部署1. 部署方式FastAPI 服务化推荐from fastapi import FastAPI import asyncio app FastAPI() workflow AIAssistantWorkflow() app.post(/run_workflow) async def run_workflow(): task asyncio.create_task(workflow.run()) await asyncio.sleep(1) await workflow.send_event(HumanReviewEvent(approveTrue)) return await task命令行持久化运行集成 Docker 云服务AWS/Azure/GCP2. 部署最佳实践启用检查点持久化存储到 Redis / 数据库设置超时时间避免死锁日志记录所有步骤异常捕获 重试机制七、学习总结工作流核心事件驱动 步骤编排 全局状态关键能力检查点断点续跑、人机协同手动触发、多事件等待开发流程定义事件 → 定义状态 → 编写步骤 → 手动触发 → 部署运行适用场景复杂 RAG、多 Agent 协作、审批流、数据处理流水线