Workflow

QueryEngine 把检索和生成焊死,中间插不上「先 CE 再决定要不要上网补检索」。Workflow 用类型化 Event 当边:一个 @step 吃某种 Event,吐另一种;返回 StopEvent 就结束。循环是「步骤再次发出自己能吃的 Event」,不是共享 State 上的自环边。

段末注释Workflow 是事件驱动编排;LangGraph 是共享 State + 显式边 + checkpointer。RAG 固定 DAG 用前者省事;要审批/断点续跑用后者。

StartEvent 信封传到 retrieve,再变成 RerankedEvent,最后 StopEvent;旁边对照 LangGraph 黑板(科普示意)


1. 一句话定位

维度 内容
角色 知识层侧的多步 DAG/弱环编排:步骤由 Event 类型接线
输入 → 输出 await workflow.run(**kwargs)StopEvent.result
典型调用入口 class X(Workflow)@stepStartEvent / StopEventctx.store
与 LangChain / LangGraph 不是 LCEL,也不是 StateGraph;节点内仍可调 Retriever / LLM

出现背景:DAG 框架把分支画成难读的边。Workflow 用普通 if 返回不同 Event 类型来分支;用「再 emit 上游 Event」来循环。校验器在运行前检查:有没有人生产某个步骤在吃的类型。


2. 前置依赖与环境

1
2
3
pip install -U llama-index-core llama-index-llms-ollama llama-index-embeddings-ollama
ollama pull qwen3.5:9b
ollama pull nomic-embed-text
  • run 为 async;timeout 单位秒
  • 独立包 llama-index-workflowsfrom workflows import Workflow;本系列用 llama_index.core.workflow 再导出路径,避免两套 import
  • 较新运行时共享状态走 ctx.store.set/get(旧文 ctx.set 可能失效)

3. 实现逻辑

1
2
3
4
5
6
7
1. 定义 Event 子类(字段=步骤间契约)
2. Workflow 子类里写 @step:参数 ev: 某Event → 返回 另Event
3. 第一个可运行步骤吃 StartEvent;run(query=...) 变成 StartEvent.query
4. retrieve 步:aretrieve,把 nodes 放进 RetrievedEvent
5. 需要跨步的小状态:await ctx.store.set("query", q)
6. 最后一步返回 StopEvent(result=...)
7. await w.run(...) 得到 result;要看中间事件则 handler.stream_events()

节点表(本篇示例)

步骤
retrieve StartEvent RetrievedEvent ev.query ctx.store[“query”]
synthesize RetrievedEvent StopEvent ev.nodes, store[“query”] result 文本

:由类型推断,不必 add_edge


4. 原理说明

主轴是:边 = Event 类型;状态默认不共享,要共享就显式放进 Context.store。

1
2
3
4
5
6
7
8
9
10
1. run(**kwargs) 构造 StartEvent,字段即 kwargs
2. 调度所有声明为吃 StartEvent 的 step(通常一个)
3. step 是 async 函数,返回一个 Event 或 list[Event]
4. 运行时按返回类型找到下一个 step;编译前 validate 生产者/消费者
5. ctx.store 是**这一次 run** 的 KV,不是跨进程 Saver
6. 若 step 再返回自己的输入类型,形成环;必须靠 timeout / 自写计数跳出
7. StopEvent 出现则停止,result 交给 await
8. handler = w.run(...);async for ev in handler.stream_events() 可观测中间 Event
少了 StopEvent → 校验失败或挂到 timeout
用 Event 塞超大 Node 列表可以,但不要把整个 Index 放进 Event 再序列化去「持久化」——那不是 checkpointer

@step 出现在步骤 2。输入输出类型必须可被静态图检查。动态 ctx.send_event 是进阶,本篇不用。

StartEvent 出现在步骤 1。可用 ev.queryev.get("query")

StopEvent(result=...) 出现在步骤 6。result 可以是 str / dict / 任意对象。

timeoutWorkflow(timeout=60)。RAG 三步本地模型建议 60~180。

与 LangGraph:那边 add_edge + reducer 合并;这边没有 reducer,后一个 Event 就是下一拍的全部输入。


5. 最小可运行示例

1
pip install -U llama-index-core llama-index-llms-ollama llama-index-embeddings-ollama
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
import asyncio
from typing import List

from llama_index.core import Document, Settings, VectorStoreIndex, get_response_synthesizer
from llama_index.core.node_parser import SentenceSplitter
from llama_index.core.schema import NodeWithScore
from llama_index.core.workflow import Event, StartEvent, StopEvent, Workflow, step, Context
from llama_index.embeddings.ollama import OllamaEmbedding
from llama_index.llms.ollama import Ollama

Settings.llm = Ollama(model="qwen3.5:9b", request_timeout=120.0, temperature=0)
Settings.embed_model = OllamaEmbedding(
model_name="nomic-embed-text",
base_url="http://localhost:11434",
)

index = VectorStoreIndex.from_documents(
[Document(text="定量 PCR 以小鼠肝脏 GAPDH 为内参。实验对象为 C57BL/6 小鼠。",
doc_id="paper_001")],
transformations=[SentenceSplitter(chunk_size=128, chunk_overlap=20)],
)


class RetrievedEvent(Event):
nodes: List[NodeWithScore]


class RAGFlow(Workflow):
@step
async def retrieve(self, ctx: Context, ev: StartEvent) -> RetrievedEvent:
query = ev.get("query")
await ctx.store.set("query", query)
nodes = await index.as_retriever(similarity_top_k=3).aretrieve(query)
return RetrievedEvent(nodes=nodes)

@step
async def synthesize(self, ctx: Context, ev: RetrievedEvent) -> StopEvent:
query = await ctx.store.get("query")
synth = get_response_synthesizer(response_mode="compact", llm=Settings.llm)
resp = await synth.asynthesize(query, nodes=ev.nodes)
cites = [n.node.doc_id for n in resp.source_nodes]
return StopEvent(result={"answer": str(resp), "cites": cites})


async def main():
w = RAGFlow(timeout=120, verbose=False)
# 输入
out = await w.run(query="实验对象是什么物种?")
print(out) # 输出:含 answer 与 cites
# 预期:answer 含小鼠;cites 含 paper_001

asyncio.run(main())

6. 重要配置参数

参数(API 名) 类型 / 默认值 功能说明 作用与影响 参考起点 / 常用范围 配置指导
timeout float,秒 整次 run 墙钟上限 过短误杀本地生成;过长空转 60~180 含 CE/多步时加大
verbose bool 打印步骤调度 只影响日志 调试 True 生产 False
ctx.store KV 跨 step 共享本次 run 的小状态 不写则下一步看不见 query 只放 query、计数 不要塞整个向量库
Event 字段 Pydantic 步骤间唯一数据契约 漏字段下一拍 AttributeError 显式 List[NodeWithScore] 与函数返回类型一致
run(**kwargs) kwargs→StartEvent 入口字段 名字必须和 ev.get 一致 query= 不要既用 query 又用 question
stream_events handler 上 观测中间 Event 不改变结果 对接 SSE 先 await 跑通再流

7. 适用 / 不适用

维度 适用 不适用
任务形态 固定 retrieve→rerank→合成;偶发 if 分支 人工审批、会话级 checkpoint
集成约束 单进程 async 多实例恢复——官方 durable 仍弱于 LangGraph Saver
工程阶段 把焊死的 QueryEngine 拆开插针 已有 StateGraph 团队——不必双轨编排

8. 易踩坑

  1. step 写成同步 def:调度器期望 async,卡死或报错。
  2. 忘了 StopEvent:校验失败或一直跑到 timeout。
  3. ctx.set 抄旧文:新版本用 ctx.store
  4. 把 Workflow 当 LangGraph:没有 thread_id 级 put/get_tuple,重启即失。

小结

  • 边是 Event 类型,不是 add_edge
  • StartEvent.kwargs → … → StopEvent.result
  • 跨步小状态放 ctx.store
  • 要 HITL / 持久化线程:检索步骤留下,外壳换 LangGraph。

参考链接

-------------本文结束感谢您的阅读-------------