LangGraph 高级特性学习文档
本文档基于
langgraph_high_level目录下 11 个示例文件整理,是 LangGraph 进阶用法的系统学习笔记。 每个章节对应一个示例文件,可对照源码学习。
目录
| 章节 | 知识点 | 对应文件 |
|---|---|---|
| 一 | 状态图三步法(基础回顾) | 所有文件通用 |
| 二 | RunnableConfig 配置注入 | 01_langgraph_config.py |
| 三 | thread_id 线程隔离 | 02_langgraph_config_thread_id.py |
| 四 | InMemorySaver 记忆机制 | 03_langgraph_memory_demo.py |
| 五 | 人工中断(Human-in-the-loop) | 04.langgraph_human_approval_demo.py |
| 六 | 并行执行(fan-out / fan-in) | 05_langgraph_并行_demo.py |
| 七 | 并行状态合并(Annotated + reducer) | 06_langgraph_并行_list_demo.py |
| 八 | 子图 Subgraph(两种用法) | 07 / 08_langgraph_subgraph_demo*.py |
| 九 | 同步流式输出 stream | 09_langgraph_stream_demo.py |
| 十 | 异步执行 ainvoke | 10_langgraph_异步_demo.py |
| 十一 | 异步流式输出 astream | 11_langgraph_异步流式_demo.py |
一、基础回顾:状态图三步法
LangGraph 编写任何工作流都遵循三步:
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
# 1️⃣ 创建状态:定义在节点之间流转的数据结构
class MyState(TypedDict):
question: str
answer: str
# 2️⃣ 创建节点:每个节点接收 state,返回更新后的 state
def think_node(state: MyState):
state["answer"] = "思考结果"
return state
# 3️⃣ 把节点串起来:add_edge 定义执行顺序
def build_graph():
gb = StateGraph(MyState)
gb.add_node("think", think_node) # 注册节点
gb.add_edge(START, "think") # 起点 → think
gb.add_edge("think", END) # think → 终点
return gb.compile() # 编译成可运行图
运行方式:graph.invoke({"question": "..."})
后续所有高级特性都建立在这三步之上,只是在不同环节加料。
二、RunnableConfig 配置注入
对应文件:01_langgraph_config.py
解决什么问题
节点函数里需要用到一些"运行时参数"(如用户名、前缀、API key 等),但又不希望写死在函数里或塞进 state 污染业务数据。
用法
from langgraph.types import RunnableConfig
def step1(state: State, config: RunnableConfig): # 👈 第二个参数
prefix = config.get("configurable", {}).get("prefix", "")
uname = config.get("configurable", {}).get("uname", "")
state["step1"] = f"{prefix}{state['input']}-{uname}"
return state
# 运行时注入
config = RunnableConfig(configurable={
"prefix": "[前缀] ",
"uname": "wangliangwu",
"renshe": "好人卡"
})
graph.invoke({"input": "你好"}, config=config)
关键点
- 节点函数签名:
def node(state, config),第二个参数固定为 RunnableConfig - 所有自定义参数都放在
config["configurable"]字典里 - 同一个 config 在整个图的每个节点都能读到(全局共享)
典型应用
- 多租户场景:传
tenant_id、user_id - 多语言场景:传
language - 调试场景:传
trace_id做链路追踪
三、thread_id 线程隔离
对应文件:02_langgraph_config_thread_id.py
解决什么问题
多个用户同时使用同一个图,怎么区分"这是 A 的会话"和"这是 B 的会话"?
用法
config = RunnableConfig(configurable={
"thread_id": "12345" # 👈 线程唯一标识
})
graph.invoke({"input": "你好"}, config=config)
关键点
thread_id是 RunnableConfig 中约定的特殊字段,LangGraph 内部会识别它- 配合记忆模块(下一章)使用才有意义
- 同一个 thread_id = 同一个会话/对话
- 不同 thread_id 之间相互隔离,互不干扰
典型应用
- Web 系统中每个用户一个 thread_id(如
session_id) - 同一用户的不同对话各开一个 thread_id
四、InMemorySaver 记忆机制
对应文件:03_langgraph_memory_demo.py
解决什么问题
默认情况下,graph.invoke() 每次都是"从零开始",前一次调用的 state 不会保留。但多轮对话场景需要"记住上次到哪了"。
用法
from langgraph.checkpoint.memory import InMemorySaver
def build_graph():
gb = StateGraph(State)
gb.add_node("add", add_node)
gb.add_edge(START, "add")
gb.add_edge("add", END)
memory = InMemorySaver() # 👈 创建内存记忆器
return gb.compile(memory) # 👈 编译时挂载
# 等价写法:gb.compile(checkpointer=memory)
graph = build_graph()
# 同一个 thread_id 多次调用 → state 会累加保留
config = RunnableConfig(configurable={"thread_id": "12345"})
graph.invoke({"new_value": 5}, config=config) # total = 5
graph.invoke({"new_value": 7}, config=config) # total = 5+7 = 12
# 换一个 thread_id → 全新开始
config2 = RunnableConfig(configurable={"thread_id": "99999"})
graph.invoke({"new_value": 10}, config=config2) # total = 10
关键点
InMemorySaver把每次调用后的 state 保存到内存,按thread_id索引- 同
thread_id再次调用时,会自动加载上次的 state 作为起点 - 不同
thread_id完全独立 - ⚠️ 内存存储重启即丢失,生产环境需换
SqliteSaver/PostgresSaver
与 thread_id 的关系
| 概念 | 作用 | 类比 |
|---|---|---|
InMemorySaver |
存储介质(在哪存) | 硬盘 |
thread_id |
数据分区键(存哪份) | 文件名 |
两者缺一不可:只有 memory 没有 thread_id,所有调用混在一起;只有 thread_id 没有 memory,每次都是空的。
五、人工中断(Human-in-the-loop)
对应文件:04.langgraph_human_approval_demo.py
解决什么问题
某些关键步骤(如审批、敏感操作、确认)需要人工介入确认后才能继续,不能让 LLM 一路跑到底。
用法
from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import InMemorySaver
def human_approval(state: State):
# 👇 触发中断:图会暂停在这里,等待外部 resume
approval = interrupt("请输入审批结果(批准 / 不批准):")
state["approval"] = approval
return state
# 必须挂载 memory,否则中断后无法恢复
memory = InMemorySaver()
graph = gb.compile(checkpointer=memory)
config = RunnableConfig(configurable={"thread_id": "12345"})
# === 第一次调用:执行到 interrupt 处自动暂停 ===
state = graph.invoke({}, config=config)
# 此时 state 包含中断信息,图暂停在 human_approval 节点
# === 第二次调用:传入人工输入,恢复执行 ===
user_decision = input("是否批准?") # 模拟人工输入
state = graph.invoke(Command(resume=user_decision), config=config)
# 👆 用 Command(resume=...) 把人工结果回灌进图,继续往后跑
关键点
interrupt(msg)是个特殊函数,调用它会暂停整个图- 暂停后必须用同一个 thread_id 才能 resume
- 恢复时用
Command(resume=值)把人工输入传回图 - ⚠️ 必须有 checkpointer,否则暂停后状态丢失,无法恢复
- 适合场景:审批流、危险操作确认、人工校对纠错
流程示意
invoke → ... → interrupt() ⛔ 暂停
↓ 状态存入 memory
↓ 等待人工
invoke(Command(resume=...)) → 继续 → END
六、并行执行(fan-out / fan-in)
对应文件:05_langgraph_并行_demo.py
解决什么问题
某些任务彼此独立(如"炒菜"和"热饭"),可以同时进行以节省时间,没必要串行等待。
用法:让一个节点连出多条边
def build_graph():
gb = StateGraph(MyState)
gb.add_node("start", start_cooking)
gb.add_node("make_vegetables", make_vegetables) # 炒菜(10秒)
gb.add_node("heat_rice", heat_rice) # 热饭(2秒)
gb.add_node("serve", serve_dish)
# 👇 fan-out:start 之后同时进入两个节点(并行)
gb.add_edge("start", "make_vegetables")
gb.add_edge("start", "heat_rice")
# 👇 fan-in:两个节点都完成后才进入 serve
gb.add_edge(["make_vegetables", "heat_rice"], "serve")
gb.add_edge("serve", END)
return gb.compile()
关键点
- 并行 = 一个节点向多个节点连边(fan-out)
- 汇总 = 多个节点向一个节点连边(fan-in),用 list 传多个源节点
- LangGraph 会等 list 中所有节点都完成,才执行下一个节点
- 并行节点之间不能写同一个 state 字段,否则会有覆盖冲突 → 见下一章
流程示意
┌→ make_vegetables ┐
start ──┤ ├→ serve → END
└→ heat_rice ───────┘
七、并行状态合并(Annotated + reducer)
对应文件:06_langgraph_并行_list_demo.py
解决什么问题
上一章的并行节点都想往 state["menu"] 这个 list 里 append 菜品,但默认情况下 state 字段会被后写的节点覆盖,导致丢数据。
用法:用 Annotated 给字段指定合并函数
from typing import Annotated
import operator
def my_merge(a, b):
"""自定义合并逻辑:这里直接返回 b(也可改成 a+b 等)"""
return b
class MyState(TypedDict):
menu: Annotated[list, my_merge] # 👈 指定合并函数
all_ready: bool
关键点
Annotated[类型, reducer函数]告诉 LangGraph:当多个节点同时写这个字段时,用 reducer 合并而不是覆盖- 内置可用的 reducer:
operator.add(list 拼接、int 累加) - 也可自定义:如
my_merge(a, b)返回b表示"后者覆盖前者" - 并行节点写同一个字段时,必须用 Annotated 指定 reducer,否则报错或丢数据
示例对照
| 写法 | 效果 |
|---|---|
menu: list |
并行写会覆盖(默认) |
menu: Annotated[list, operator.add] |
多个 list 自动拼接 |
menu: Annotated[list, my_merge] |
用自定义函数合并 |
八、子图 Subgraph
子图是把一组节点封装成"可复用的子流程",主图把它当一个节点用。有两种封装方式。
方式 1:直接把子图当节点加入主图
对应文件:07_langgraph_subgraph_demo.py
def build_subgraph():
g = StateGraph(SubState) # 👈 子图有自己的 State
g.add_node("hello", say_hello)
g.add_node("bye", say_bye)
g.add_edge(START, "hello")
g.add_edge("hello", "bye")
g.add_edge("bye", END)
return g.compile()
def build_main_graph():
subgraph = build_subgraph() # 👈 先创建子图实例
g = StateGraph(MainState)
g.add_node("start", start_main)
g.add_node("sub", subgraph) # 👈 直接把子图当节点加入
g.add_node("finish", finish_main)
g.add_edge(START, "start")
g.add_edge("start", "sub")
g.add_edge("sub", "finish")
g.add_edge("finish", END)
return g.compile()
特点:
- 写法简单,一行
add_node("sub", subgraph)搞定 - ⚠️ 主图和子图共享同名字段才能传数据(隐式传递)
- 子图内部对 main state 中没有的字段无感知
方式 2:包装节点手动调用子图
对应文件:08_langgraph_subgraph_demo2.py
def run_subgraph(state: MainState):
subgraph = build_subgraph()
# 👇 手动构造子图输入(字段名/数据可任意转换)
sub_state = {"name": "小帅"}
result = subgraph.invoke(sub_state)
# 👇 手动把子图输出合并回主图 state
state["sub_output"] = result
return state
def build_main_graph():
g = StateGraph(MainState)
g.add_node("start", start_main)
g.add_node("sub", run_subgraph) # 👈 包装函数,不是子图本身
...
特点:
- 灵活:可手动转换主图↔子图的字段映射
- 解耦:主图 State 和子图 State 完全独立,无字段共享要求
- 适合:子图来自第三方、字段名对不上、需要预处理后再调用
两种方式对比
| 维度 | 方式1(直接加入) | 方式2(包装节点) |
|---|---|---|
| 代码量 | 少 | 多 |
| 字段传递 | 隐式(同名共享) | 显式(手动构造) |
| 灵活性 | 低 | 高 |
| 耦合度 | 高 | 低 |
| 推荐场景 | 主子图同团队开发 | 子图复用/第三方 |
九、同步流式输出 stream
对应文件:09_langgraph_stream_demo.py
解决什么问题
invoke 是"跑完才一次性返回",节点很多时用户要干等。stream 让每个节点完成时立刻吐出结果,提升交互体验。
用法
# invoke:阻塞到全部完成
result = graph.invoke({"question": "..."})
# stream:每个节点完成时立刻返回一次
for chunk in graph.stream({"question": "为什么天空是蓝色的?"}):
print(chunk)
# 输出示例:
# {"思考节点": {"thought": "是怎么染色的呢?"}}
# {"回答节点": {"answer": "是用彩笔染色的。"}}
关键点
stream返回生成器,用for chunk in ...接收- 每个 chunk 是一个 dict,key 是节点名,value 是该节点返回的 state 增量
- 适合:命令行调试、本地脚本、不需要并发的场景
十、异步执行 ainvoke
对应文件:10_langgraph_异步_demo.py
解决什么问题
同步 invoke 会阻塞当前线程。Web 服务(FastAPI/Streamlit)中需要异步执行,才能同时处理多个请求。
用法
import asyncio
async def main(state: MyState):
graph = build_graph()
result = await graph.ainvoke(state) # 👈 异步版本
return result
# 启动异步事件循环
asyncio.run(main({"question": "..."}))
与 asyncio.create_task 配合实现并发
async def run():
task1 = asyncio.create_task(main({"question": "..."})) # LangGraph 任务
task2 = asyncio.create_task(fastprint()) # 其他异步任务
await task1
await task2
asyncio.run(run())
关键点
- 同步版
invoke→ 异步版ainvoke(加a前缀) - 异步函数必须在事件循环中跑,用
asyncio.run()或await - 适合:FastAPI、Web 服务、需要并发处理多请求的场景
十一、异步流式输出 astream
对应文件:11_langgraph_异步流式_demo.py
解决什么问题
结合"流式"和"异步"两者的优点:既能边算边输出,又不阻塞事件循环。是生产 Web 服务的首选方案。
用法
async def main(state: MyState):
graph = build_graph()
async for event in graph.astream(state): # 👈 异步流式
print(event)
asyncio.run(main({"question": "..."}))
关键点
astream是异步生成器,必须用async for接收- 不能用
asyncio.run(main())之外的同步 for 循环 - 适合:FastAPI SSE 流式响应、实时聊天界面
三种调用方式对照表
| 方法 | 同步/异步 | 是否流式 | 适用场景 |
|---|---|---|---|
graph.invoke(input) |
同步 | 否 | 命令行脚本、简单测试 |
graph.stream(input) |
同步 | 是 | 本地调试、看每步输出 |
graph.ainvoke(input) |
异步 | 否 | Web 服务单次返回 |
graph.astream(input) |
异步 | 是 | Web 服务流式响应(推荐) |
十二、知识点速查表
12.1 核心 API 速查
| API | 作用 | 出处 |
|---|---|---|
StateGraph(State) |
创建图构造器 | 所有文件 |
gb.add_node(name, fn) |
注册节点 | 所有文件 |
gb.add_edge(src, dst) |
加边(顺序执行) | 所有文件 |
gb.add_edge([a, b], c) |
多源汇总(fan-in) | 并行示例 |
gb.compile(checkpointer=...) |
编译图,挂载记忆器 | memory/审批示例 |
graph.invoke(input, config) |
同步执行 | 所有文件 |
graph.stream(input, config) |
同步流式 | 09 |
graph.ainvoke(input, config) |
异步执行 | 10 |
graph.astream(input, config) |
异步流式 | 11 |
interrupt(msg) |
触发人工中断 | 04 |
Command(resume=...) |
恢复中断执行 | 04 |
12.2 关键导入
from langgraph.graph import StateGraph, START, END
from langgraph.types import RunnableConfig, interrupt, Command
from langgraph.checkpoint.memory import InMemorySaver
from typing import TypedDict, Annotated
12.3 状态定义对比
# 基础状态
class State(TypedDict):
field: str
# 带合并函数的状态(并行场景必备)
class State(TypedDict):
menu: Annotated[list, operator.add] # 多节点写同一字段时自动合并
12.4 节点函数签名
# 最简签名
def node(state): ...
# 带配置(需要读 config 时)
def node(state, config: RunnableConfig): ...
# 异步节点
async def node(state): ...
十三、学习路径建议
按以下顺序学习,每个文件都跑一遍代码看输出:
01 config → 学会往节点传运行时参数
↓
02 thread_id → 理解会话隔离的概念
↓
03 memory → 掌握多轮记忆(thread_id + memory 配合)
↓
04 human_approval → 学会中断恢复(依赖 memory)
↓
05 并行 → 学会 fan-out / fan-in
↓
06 并行 list → 解决并行写冲突(依赖 05)
↓
07/08 子图 → 学会复用子流程
↓
09 stream → 同步流式
↓
10 ainvoke → 异步执行
↓
11 astream → 异步流式(生产推荐)
学习要点
- 每章对应一个独立 demo,可单独运行(都有
if __name__ == "__main__") - memory + thread_id 是核心:4/5/6/10/11 章都依赖它
- 生产环境三件套:
astream+InMemorySaver+thread_id(或换持久化存储) - 本项目主流程参考:
__004__langgraph_more_nodes/就是这些特性的综合应用
十四、与本项目主流程的对应关系
| 本项目主流程 | 用到的 LangGraph 特性 | 对应 demo |
|---|---|---|
full_workflow.py 多节点串联 |
StateGraph + add_edge | 所有基础 |
match_entity_node 等 9 节点 |
节点函数 + State 流转 | 01 |
| FastAPI 多用户隔离 | thread_id + InMemorySaver | 02 + 03 |
check_cypher_node 重试回退 |
条件边(add_conditional_edges) | 05 思想 |
/process_stream 流式接口 |
astream + asyncio.Queue | 11 |
| 多轮对话 | InMemorySaver 持久化 history | 03 |
学完这 11 个 demo,再回头看
__004__langgraph_more_nodes/full_workflow.py会非常轻松。