作为LangGraph节点

检索质量已经交给 LlamaIndex:切块、混合召回、source_nodes 都能审计。真正卡住产品的是:工具失败要再问一次、发信前要人点头、进程重启后续跑。这些是 LangGraphStatecheckpointerinterrupt。做法是把 QueryEngine 当作图上一个节点(或节点内的一次调用),而不是用 Workflow 再实现一遍 Pregel。

段末注释边界 = LlamaIndex 产出 Response + source_nodes;LangGraph 决定何时调用、如何写入 messages、如何按 thread_id 存取。

右轨:LangGraph 环上的 rag_node 内嵌 QueryEngine,checkpoint 车票 thread_id(科普示意)


1. 一句话定位

维度 内容
角色 把知识层 QueryEngine 嵌进编排层 StateGraph
输入 → 输出 State["messages"] 最后一条用户话 → qe.query → 追加 AIMessage(含引用)
典型调用入口 节点函数内 qe.query / await qe.aquerygraph.compile(checkpointer=...)
与纯 LlamaIndex 不再用 FunctionAgent/Workflow 管环;它们只适合短 DAG / 短工具环原型

出现背景:社区常见「LlamaIndex 管知识、LangGraph 管控制流」。LangChain 生态也提供 Retriever 适配器;本篇用节点内直接调 QE,少一层包装,引用字段不丢失。


2. 前置依赖与环境

1
2
3
pip install -U langgraph langchain-core llama-index-core llama-index-llms-ollama llama-index-embeddings-ollama
ollama pull qwen3.5:9b
ollama pull nomic-embed-text
  • Python 3.10+;LangGraph 与 langchain-core 版本需匹配(本仓库 LangGraph 目录锚点 1.2.x)
  • 图节点若跑在 FastAPI 里:用 async 图 + qe.aquery,不要同步 query 堵事件循环
  • 本篇示例无 HITL;interrupt 用法见 LangGraph 专篇,节点名填 rag 即可

3. 实现逻辑

节点表

节点名 职责
rag 调 LlamaIndex QueryEngine messages[-1].content messages 追加 AI 回复

边表

目标 类型
START rag 固定
rag END 固定

需要「答案不行再检索」时,加条件边回到 rag 或到 tools 节点——环画在图上,不要在 qe.query 内部 while。

1
2
3
4
5
6
7
1. 进程启动时建好 Index / QueryEngine(只一次,不要每个节点重建)
2. State: messages: Annotated[list[BaseMessage], add_messages]
3. rag_node(state):q = 最后一条 Human/用户内容
4. resp = qe.query(q)
5. 把 str(resp) 与 source_nodes 的 doc_id/section 拼进 AIMessage.content
6. return {"messages": [AIMessage(...)]} # 部分更新,靠 add_messages 追加
7. compile(checkpointer=InMemorySaver());invoke(..., config={"configurable":{"thread_id":"t1"}})

字段级变形

1
2
3
HumanMessage("实验对象是什么物种?")
→ qe.query → Response(response="小鼠……", source_nodes=[NodeWithScore(doc_id=paper_001)])
→ AIMessage(content="小鼠……\n[cite] paper_001 methods")

4. 原理说明

主轴是:图调度 LangGraph;一次知识访问 LlamaIndex;两套状态不要混写。

1
2
3
4
5
6
7
8
9
10
1. graph.invoke 按 thread_id 读 checkpoint(若有 Saver)
2. 进入 rag 节点,拿到合并后的 state["messages"]
3. 只把「当前用户问题」交给 qe.query——不要把整段 LangChain 历史塞进 LlamaIndex,除非你有意做对话检索
4. QE 内部走 retrieve → postprocess → synthesize,与图无关
5. 节点返回 partial update;add_messages 追加,而不是覆盖整个 messages
6. superstep 结束 Saver.put;到 END 返回
7. 同一 thread_id 第二次 invoke:历史还在图的 messages 里;QE 侧默认仍是无状态 DAG
少了第 5 步 reducer → 第二节点写 messages 会盖掉 Human
少了第 3 步把 20 轮聊天当 query → 检索词被稀释
在第 4 步里写 while tool:检查点看不到环,HITL 插不进去

qe.query / aquery 出现在步骤 4。返回 Response。引用必须在节点里拷到 AIMessage,因为 LangGraph 不会自动懂 source_nodes

add_messages 出现在步骤 5。与 LangGraph 状态专篇相同。

InMemorySaver 出现在步骤 7。开发用;生产换 Postgres/Sqlite Saver。LlamaIndex Context 不能替代 Saver。

可选:把 QE 包成 LangChain Tool,再 bind_tools 走完整 ReAct 图。本篇故意不这么做,以免两套 Document/Message 转换丢 metadata。


5. 最小可运行示例

1
pip install -U langgraph langchain-core 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
53
from typing import Annotated, TypedDict

from langchain_core.messages import AIMessage, BaseMessage, HumanMessage
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.graph.message import add_messages
from llama_index.core import Document, Settings, VectorStoreIndex
from llama_index.core.node_parser import SentenceSplitter
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", metadata={"section": "methods"})],
transformations=[SentenceSplitter(chunk_size=128, chunk_overlap=20)],
)
qe = index.as_query_engine(similarity_top_k=3)


class State(TypedDict):
messages: Annotated[list[BaseMessage], add_messages]


def rag_node(state: State) -> dict:
"""读最后一条用户消息,调 QueryEngine,把答案与 cite 追加为 AIMessage。"""
q = state["messages"][-1].content
resp = qe.query(q)
cites = ", ".join(
f"{n.node.doc_id}:{n.node.metadata.get('section', '')}"
for n in resp.source_nodes
)
return {"messages": [AIMessage(content=f"{resp}\n[cite] {cites}")]}


builder = StateGraph(State)
builder.add_node("rag", rag_node)
builder.add_edge(START, "rag")
builder.add_edge("rag", END)
graph = builder.compile(checkpointer=InMemorySaver())

config = {"configurable": {"thread_id": "demo-li"}}
out = graph.invoke(
{"messages": [HumanMessage(content="实验对象是什么物种?")]},
config,
)
print(out["messages"][-1].content)
# 预期:含小鼠,以及 [cite] paper_001:methods

需要人工批准再生成:compile(interrupt_before=["rag"]),批准后 Command(resume=...)——API 细节在 LangGraph HITL 专篇。


6. 重要配置参数

参数(API 名) 类型 / 默认值 功能说明 作用与影响 参考起点 / 常用范围 配置指导
节点返回 messages list[BaseMessage] 只返回增量消息 返回全量且无 reducer 会覆盖 [AIMessage(...)] 必须配 add_messages
thread_id str,configurable Saver 第一层键 换 ID = 新会话 user-{id}-conv 多租户唯一
checkpointer Saver 图状态持久化 InMemory 进程退出即失 开发内存 / 生产 PG 与 LlamaIndex persist 目录无关
interrupt_before list[节点名] 进节点前暂停 不配则无法 HITL 高风险节点名 检索节点一般不打断,生成/发信才打断
qe 生命周期 进程级单例 避免每次 invoke 重建索引 建在节点函数里会每次 embedding 模块加载时建好 persist 索引与 Saver 分开
aquery async 异步图里的知识调用 同步 query 堵 event loop FastAPI 必 async ainvoke 成对

7. 适用 / 不适用

维度 适用 不适用
任务形态 要引用的 RAG + 要环/审批/恢复 单次问答——不用上图
集成约束 已同时用两套库 团队只会 LlamaIndex Workflow——固定 DAG 不必硬上 LangGraph
工程阶段 服务化 Agent 用 LangChain Tool 适配器却丢掉 metadata——先直连 QE

8. 易踩坑

  1. 每个节点 from_documents:延迟爆炸。
  2. messages 无 reducer
  3. 把 整段对话当 QE query
  4. 以为 LlamaIndex persist 能恢复 LangGraph 线程:两套存储,互不相认。
  5. 在 rag_node 里 while tool_calls:检查点看不到中间态。

小结

  • QE 当节点,图管环。
  • 引用在节点里写进 AIMessage,LangGraph 不认识 source_nodes
  • Saver + thread_id 管会话;Index persist 管向量。
  • HITL 用 interrupt_*,不要用 FunctionAgent Context 冒充。

参考链接

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