LlamaIndex Workflows 入门:实用指南

open-source入门14 分钟阅读2026/7/11

上个月,我正在为一个客户搭建文档处理流水线。他们需要从保险理赔单中提取数据,与保单文档进行交叉比对,并标记潜在的欺诈风险。我一开始用了最直接的方式——串行调用 LLM:先提取,再验证,最后评估。在一切顺利的“理想路径”下,这套流程跑得挺好。但当我需要针对缺失的信息进行循环补全,或者让不同类型的文档走不同的验证步骤时,我原本精心构建的调用链就变成了一团乱麻,里面塞满了条件逻辑和嵌套回调。

就在这时候,我偶然发现了 LlamaIndex Workflows。传统的调用链强迫你使用死板的 A→B→C 执行模型,而 Workflows 则不同,它采用了事件驱动架构。各个步骤(Step)会发出事件,其他步骤则监听这些事件并做出反应。需要循环?只需让一个步骤发出一个事件,去触发前面的步骤即可。需要条件分支?根据你的逻辑发出不同的事件就行了。这是一个小小的概念转变,却彻底改变了你构建复杂应用的方式。

接下来,我将带你了解我是如何实际上手的,包括那些让我踩坑的地方。

核心概念:事件与步骤

在写代码之前,你需要先搞懂两个基础构建块:

  1. 事件 是数据容器。它们负责在步骤之间传递信息。
  2. 步骤 是用 @step 装饰器修饰的函数,它们监听特定的事件,并发出其他事件。

你可以把它想象成一个无线电广播系统。一个步骤调到特定的频率(事件类型),处理听到的内容,然后再用另一个频率广播出去。多个步骤可以监听同一个事件,一个步骤也可以发出多种不同的事件。

环境配置

首先,安装核心包。这里我使用的是 OpenAI,但你也可以替换成任何其他模型提供商:

pip install "llama-index-core>=0.11.16" llama-index-utils-workflow
pip install llama-index-llms-openai

你还需要配置好 API 密钥:

import os
os.environ["OPENAI_API_KEY"] = "your-key-here"

构建一个带重排序的实用 RAG 工作流

咱们来点实际的:构建一个 RAG 流水线,让它检索文档、对相关性进行重排序,然后生成答案。我选这个例子是因为大家都懂 RAG,而且它已经能从事件驱动模型中获益了——尤其是当你后续想加入重试逻辑或条件路由时。

第一步:定义你的事件

事件其实就是 Pydantic 模型。这其实是第一个让我困惑的地方——我一开始总想在步骤之间直接传原始字典。你必须定义明确的事件类:

from llama_index.core.workflow import StartEvent, StopEvent, Event

class RetrieveEvent(Event):
    query: str

class RerankEvent(Event):
    query: str
    documents: list

class SynthesizeEvent(Event):
    query: str
    ranked_documents: list

StartEventStopEvent 是内置的。当你调用 workflow.run() 时,StartEvent 负责启动一切;而 StopEvent 负责终止工作流并返回结果。

第二步:定义工作流类和步骤

现在我们来创建工作流类并添加步骤。每个步骤通过类型提示来声明它监听哪个事件,以及可以发出哪些事件:

from llama_index.core.workflow import Workflow, step
from llama_index.core import VectorStoreIndex, SimpleDirectoryReader, Settings
from llama_index.core.postprocessor import SentenceTransformerRerank

class RAGWorkflow(Workflow):
    
    @step
    async def retrieve(self, ev: StartEvent) -> RetrieveEvent:
        query = ev.get("query", "")
        if not query:
            raise ValueError("Query is required!")
        
        # 加载并索引文档(在生产环境中,你会预先构建好索引)
        documents = SimpleDirectoryReader("./data").load_data()
        index = VectorStoreIndex.from_documents(documents)
        retriever = index.as_retriever(similarity_top_k=10)
        
        nodes = retriever.retrieve(query)
        documents = [node.node.get_content() for node in nodes]
        
        return RetrieveEvent(query=query, documents=documents)
    
    @step
    async def rerank(self, ev: RetrieveEvent) -> RerankEvent:
        # 使用重排序器来缩小结果范围
        reranker = SentenceTransformerRerank(
            model="cross-encoder/ms-marco-MiniLM-L-2-v2",
            top_n=3
        )
        
        # 在真实的实现中,你会传递实际的 NodeWithScore 对象
        # 为了简单起见,我们这里只把排名靠前的文档传下去
        ranked_docs = ev.documents[:3]  # 为演示作简化
        
        return RerankEvent(query=ev.query, ranked_documents=ranked_docs)
    
    @step
    async def synthesize(self, ev: RerankEvent) -> StopEvent:
        from llama_index.llms.openai import OpenAI
        llm = OpenAI(model="gpt-4o")
        
        context = "\n\n---\n\n".join(ev.ranked_documents)
        prompt = f"""Based on the following context, answer the question.
        
Context:
{context}

Question: {ev.query}

Answer:"""
        
        response = llm.complete(prompt)
        return StopEvent(result=str(response))

第三步:运行工作流

async def main():
    workflow = RAGWorkflow(timeout=60, verbose=True)
    result = await workflow.run(query="What are the key risk factors mentioned in the documents?")
    print(result)

# 在 Jupyter notebook 中运行:
import asyncio
asyncio.run(main())

在开发阶段,verbose=True 这个参数简直太有用了。它会打印出可视化的追踪记录,显示哪个事件触发了哪个步骤,这样你就能实时看到整个流程的执行情况。

我踩过的坑(以及经验教训)

坑 1:忘记声明发出的事件类型。 @step 装饰器会利用类型提示来构建你的工作流图。如果你的步骤有条件地发出不同的事件,你需要使用联合类型:

@step
async def validate(self, ev: RetrieveEvent) -> RerankEvent | StopEvent:
    if not ev.documents:
        return StopEvent(result="No documents found")
    return RerankEvent(query=ev.query, documents=ev.documents)

我一开始没这么做,结果报了一些让人摸不着头脑的错误,说找不到事件处理器。框架需要在解析时就知道所有可能发出的事件,这样才能验证工作流图。

坑 2:试图用全局变量在步骤间共享状态。 Workflows 的设计理念是步骤之间无状态。如果你需要在各个步骤间持久化数据,请使用 Context 对象:

class MyWorkflow(Workflow):
    
    @step
    async def step_one(self, ev: StartEvent, ctx: Context) -> RetrieveEvent:
        await ctx.set("original_query", ev.get("query"))
        return RetrieveEvent(query=ev.get("query"))
    
    @step
    async def step_two(self, ev: RerankEvent, ctx: Context) -> StopEvent:
        # 获取共享状态
        original_query = await ctx.get("original_query")
        return StopEvent(result=f"Processed: {original_query}")

注意 ctx: Context 这个参数。框架会自动注入它。对于任何稍微复杂点的工作流,这个模式都是必不可少的。

坑 3:没有正确处理异步问题。 Workflows 本质上是异步的。如果你在 Jupyter notebook 中运行,可能需要用到 nest_asyncio 来避免事件循环冲突:

import nest_asyncio
nest_asyncio.apply()

加入循环:真正的威力所在

这才是 Workflows 真正完胜线性链的地方。假设我希望我的 RAG 流水线在初始检索结果看起来不相关时能够重试:

class RetryEvent(Event):
    query: str
    attempt: int

class RAGWorkflowWithRetry(Workflow):
    
    @step
    async def retrieve(self, ev: StartEvent | RetryEvent, ctx: Context) -> RetrieveEvent | StopEvent:
        query = ev.get("query", "") if isinstance(ev, StartEvent) else ev.query
        attempt = await ctx.get("attempt", default=0) + 1
        await ctx.set("attempt", attempt)
        
        if attempt > 3:
            return StopEvent(result="Could not find relevant documents after 3 attempts")
        
        # ... 检索逻辑 ...
        nodes = retriever.retrieve(query)
        
        # 检查相关性 - 如果很差,就循环回去
        if self._is_poor_retrieval(nodes):
            return RetryEvent(query=query, attempt=attempt)
        
        return RetrieveEvent(query=query, documents=[n.node.get_content() for n in nodes])
    
    def _is_poor_retrieval(self, nodes):
        # 简单的启发式方法:检查最高得分是否低于阈值
        return all(n.score < 0.5 for n in nodes) if nodes else True

你试试用线性链把这种循环写得干净点试试。最后你肯定会搞出一堆 while 循环、让人头疼的状态管理,以及很容易出错的流程控制。但在 Workflows 里,只需要发出一个不同类型的事件就搞定了。

可视化你的工作流

有一个功能我后来才体会到它的好:你可以为工作流生成可视化的图表:

from llama_index.utils.workflow import draw_all_possible_flows

draw_all_possible_flows(RAGWorkflowWithRetry, filename="workflow.html")

这会生成一个 HTML 文件,里面有一张交互式图表,展示了工作流中所有可能的路径。有一次我遇到一个 Bug,某个步骤怎么都触发不了,结果这张图表立刻让我发现,我忘了在联合返回类型中包含一个事件类型。这让我省了一个小时的调试时间。

实用建议与诚实的局限性

建议:

  • 开发时一定要开启 verbose=True。事件追踪是你最好的调试工具。
  • 保持步骤短小精悍。一个步骤应该只做好一件事,就像写好函数一样。
  • 需要在步骤间持久化任何数据时,请使用 Context——别想着用实例变量来偷懒。
  • 给事件起个描述性强的名字。当你阅读追踪输出时,你会感谢自己的。

我碰到的局限性:

  • 事件驱动模型的学习曲线确实存在,尤其是如果你习惯了命令式的链式调用。在豁然开朗之前,你大概会有一两天的时间在嘀咕:“这要是写个普通函数不是更简单吗”。
  • 当步骤发出多种事件类型时,调试可能会比较棘手。可视化功能能帮上忙,但在复杂的分支逻辑中理清调用链仍然需要十分细心。
  • 文档虽然在不断完善,但在扇出/扇入(多个步骤并行处理同一个事件)等高级模式方面仍有欠缺。我不得不去阅读源码才搞明白其中一些模式。
  • 对于非常简单的流水线(真的是那种没有任何分支的 A→B→C),相比于基础链,Workflows 会增加一些额外开销。在需要灵活性时才去用它,别把它当默认选项。

LlamaIndex Workflows 优雅地解决了我处理文档的问题。那个原本正变成一堆难以维护的条件判断的保险理赔流水线,现在变成了一组干净的、自带说明的步骤,数据流向清晰明了。事件驱动模型不仅仅是一种不同的语法——它真正改变了你构思和构建复杂 LLM 应用的方式,而一旦你转过弯来,你就再也不想回到线性链了。

相关 Agent

M

Meta AI

Meta AI 是一个开源AI平台,用于研究和开发先进的语言模型及生成式AI工具。

了解更多 →