LangGraph高级特性学习文档

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_iduser_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         → 异步流式(生产推荐)

学习要点

  1. 每章对应一个独立 demo,可单独运行(都有 if __name__ == "__main__"
  2. memory + thread_id 是核心:4/5/6/10/11 章都依赖它
  3. 生产环境三件套astream + InMemorySaver + thread_id(或换持久化存储)
  4. 本项目主流程参考__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 会非常轻松。