LangGraph 作为 LangChain 生态中用于构建复杂多智能体工作流的核心框架,其高级特性赋予了开发者灵活构建生产级工作流的能力。本文将深度拆解流式处理 (Streaming)、状态持久化 (Persistence)、时间回溯 (Time-Travel)、子图 (Subgraphs) 四大核心特性,结合实战代码解析底层逻辑与最佳实践,助力开发者掌握工业级工作流构建技巧。
一、流式处理 (Streaming):实时感知工作流状态
核心价值
传统的 invoke 方法只能等待工作流全量执行完成后返回结果,而流式处理支持分阶段获取状态更新,适用于实时展示进度、LLM 令牌逐字输出、自定义数据推送等场景,大幅提升交互体验。
核心模式与使用场景
LangGraph 提供多类流式模式,覆盖不同的实时数据获取需求:
| updates | 仅输出每一步的增量状态(节点返回的更新数据) | 实时展示工作流执行步骤(如「步骤 1 完成」「步骤 2 执行中」) |
| values | 输出每一步后的完整状态(包含所有历史数据) | 监控工作流全量状态,便于问题排查 |
| custom | 支持节点内推送自定义业务数据 | 推送进度百分比、业务日志、用户提示等非状态核心数据 |
| messages | 逐 Token 输出 LLM 响应(元组:(token片段, 元数据)) | 模拟 ChatGPT 打字效果,提升用户体验 |
| 多模式组合 | 如 ["values", "custom"],同时获取多维度数据 | 复杂场景下的全维度监控(状态 + 自定义进度) |
核心代码示例
1. 基础流式(updates/values)
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
# 定义状态
class AtguiguState(TypedDict):
question: str
answer: str
steps: list
# 节点函数
def think(state: AtguiguState) -> AtguiguState:
return {"steps": [f"分析问题: {state['question']}"]}
def respond(state: AtguiguState) -> AtguiguState:
return {"answer": f"回答: {state['question']}"}
# 构建图并流式执行
graph = (
StateGraph(AtguiguState)
.add_node("think", think)
.add_node("respond", respond)
.add_edge(START, "think")
.add_edge("think", "respond")
.add_edge("respond", END)
.compile()
)
# updates模式:仅看增量
for chunk in graph.stream({"question": "今天天气怎么样?"}, stream_mode="updates"):
print("增量更新:", chunk)
# values模式:看全量状态
for chunk in graph.stream({"question": "今天天气怎么样?"}, stream_mode="values"):
print("全量状态:", chunk)
2. 自定义流式数据
from langgraph.config import get_stream_writer
def node_with_custom_streaming(state: State) -> State:
# 获取流写入器
writer = get_stream_writer()
# 推送自定义进度数据
writer({"progress": "步骤1: 分析查询", "status": "running"})
writer({"progress": "步骤2: 生成结果", "status": "completed"})
return {"answer": f"处理结果: {state['query']}"}
# 执行自定义流式
for chunk in graph.stream({"query": "hello world"}, stream_mode="custom"):
print("自定义数据:", chunk)
3. LLM 逐 Token 流式
def llm_node(state: State):
# 初始化大模型
llm = init_chat_model(
model="qwen-plus",
model_provider="openai",
api_key=os.getenv("aliQwen-api"),
base_url="https://dashscope.aliyuncs.com/compatible-mode/v1"
)
# 流式调用LLM并输出
llm_result = llm.invoke([("user", state["query"])])
return {"answer": llm_result}
# messages模式:逐Token输出
for chunk, meta_data in graph.stream({"query": "小学生作文:我的一天"}, stream_mode="messages"):
print(chunk.content, end="") # 模拟打字效果
关键注意事项
- custom 模式需通过 get_stream_writer() 主动推送数据,否则无输出;
- 多模式组合时需以列表形式传入(如 stream_mode=["values", "custom"]);
- messages 模式仅对调用 LLM 的节点生效,需确保节点内有 LLM 调用逻辑。
二、状态持久化 (Persistence):工作流状态的可靠存储
核心价值
默认情况下 LangGraph 状态仅存于内存,程序重启后丢失。状态持久化支持将工作流的检查点 (Checkpoint) 存储到外部介质,实现:
- 多轮对话上下文保留;
- 工作流中断后恢复执行;
- 生产环境下的状态可追溯、可审计。
核心存储方案
LangGraph 提供标准化的 BaseCheckpointSaver 接口,适配不同存储介质:
| 内存存储 | InMemorySaver | 本地测试、临时验证 | 无需配置,程序关闭后数据丢失 |
| SQLite | SqliteSaver | 单机生产、轻量部署 | 本地文件存储,无需数据库服务 |
| PostgreSQL | PostgresSaver | 分布式生产、高可用 | 支持集群,工业级存储 |
| Redis | RedisSaver | 高性能缓存、临时状态 | 内存 + 磁盘,读写速度快 |
核心代码示例
1. 内存持久化(测试用)
from langgraph.checkpoint.memory import InMemorySaver
from typing import Annotated
import operator
# 定义带持久化的状态
class PersistenceDemoState(TypedDict):
messages: Annotated[list, operator.add]
step_count: Annotated[int, operator.add]
# 节点函数
def step_one(state: PersistenceDemoState) -> dict:
return {"messages": ["执行步骤1"], "step_count": 1}
# 构建图并配置持久化
builder = StateGraph(PersistenceDemoState)
builder.add_node("step_one", step_one)
builder.add_edge(START, "step_one")
builder.add_edge("step_one", END)
# 配置内存检查点
graph = builder.compile(checkpointer=InMemorySaver())
# 执行并指定线程ID(会话标识)
config = {"configurable": {"thread_id": "user_123"}}
# 首次执行
graph.invoke({"messages": ["开始"], "step_count": 0}, config)
# 恢复执行(直接从检查点获取状态)
graph.invoke(None, config)
# 查看历史状态
history = graph.get_state_history(config)
for checkpoint in history:
print("历史状态:", checkpoint.values)
2. SQLite 持久化(生产轻量版)
import sqlite3
from langgraph.checkpoint.sqlite import SqliteSaver
# 连接SQLite数据库
conn = sqlite3.connect("workflow_state.db", check_same_thread=False)
sqlite_saver = SqliteSaver(conn=conn)
# 编译图并使用SQLite存储
graph = builder.compile(checkpointer=sqlite_saver)
# 执行并存储状态
config = {"configurable": {"thread_id": "user_456"}}
graph.invoke({"messages": []}, config)
# 关闭连接
conn.close()
关键注意事项
- 每个会话需指定唯一 thread_id,用于区分不同用户 / 会话的状态;
- 生产环境优先使用 PostgreSQL/Redis,避免 SQLite 单机瓶颈;
- operator.add 是状态合并策略,确保多次执行的状态增量累加(而非覆盖)。
三、时间回溯 (Time-Travel):工作流的「时光机」
核心价值
时间回溯是状态持久化的高阶应用,支持:
- 回溯到工作流任意历史检查点;
- 修改历史状态后重新执行(如修正错误的节点输出);
- 对比不同执行路径的结果,实现「分支实验」。
核心流程
核心代码示例
import uuid
from langgraph.checkpoint.memory import InMemorySaver
# 定义故事生成状态
class StoryState(TypedDict):
character: str # 角色
setting: str # 场景
plot: str # 剧情
# 节点函数:生成角色、场景、剧情
def create_character(state: StoryState):
return {"character": "一只会说话的猫"}
def set_setting(state: StoryState):
return {"setting": "神秘图书馆"}
def develop_plot(state: StoryState):
return {"plot": f"{state['character']}在{state['setting']}发现发光的书"}
# 构建工作流
workflow = StateGraph(StoryState)
workflow.add_node("create_character", create_character)
workflow.add_node("set_setting", set_setting)
workflow.add_node("develop_plot", develop_plot)
workflow.add_edge(START, "create_character")
workflow.add_edge("create_character", "set_setting")
workflow.add_edge("set_setting", "develop_plot")
workflow.add_edge("develop_plot", END)
# 配置持久化
graph = workflow.compile(checkpointer=InMemorySaver())
# 1. 首次执行生成故事
config1 = {"configurable": {"thread_id": str(uuid.uuid4())}}
story1 = graph.invoke({}, config1)
print("原故事:", story1)
# 2. 获取历史检查点(找到角色生成后的状态)
history = list(graph.get_state_history(config1))
character_checkpoint = history[2] # 索引2对应create_character执行后
# 3. 修改历史状态(把猫改成龙)
new_config = graph.update_state(
character_checkpoint.config,
values={"character": "一只会飞的龙"}
)
# 4. 从修改后的状态恢复执行
story2 = graph.invoke(None, new_config)
print("修改后的故事:", story2)
关键注意事项
- 检查点索引需通过 get_state_history 确认,不同节点执行后索引不同;
- update_state 仅修改指定检查点的状态,不影响其他历史记录;
- 时间回溯依赖持久化,需提前配置 CheckpointSaver。
四、子图 (Subgraphs):工作流的模块化复用
核心价值
子图允许将一组节点封装为独立的「子工作流」,并作为单个节点嵌入到主图中,实现:
- 模块化开发(如多智能体拆分);
- 工作流复用(相同逻辑无需重复编写);
- 复杂工作流的分层管理(主图管流程,子图管细节)。
核心规则
核心代码示例
1. 基础子图嵌入
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from operator import add
# 定义共享状态
class AtguiguState(TypedDict):
messages: Annotated[list[str], add]
# 子图节点函数
def sub_node(state: AtguiguState) -> AtguiguState:
return {"messages": ["子图响应"]}
# 构建子图
subgraph_builder = StateGraph(AtguiguState)
subgraph_builder.add_node("sub_node", sub_node)
subgraph_builder.add_edge(START, "sub_node")
subgraph_builder.add_edge("sub_node", END)
subgraph = subgraph_builder.compile()
# 构建主图(嵌入子图)
builder = StateGraph(AtguiguState)
builder.add_node("subgraph_node", subgraph) # 子图作为单个节点
builder.add_edge(START, "subgraph_node")
builder.add_edge("subgraph_node", END)
graph = builder.compile()
# 执行主图(触发子图执行)
result = graph.invoke({"messages": ["主图初始消息"]})
print("最终状态:", result) # 输出: ["主图初始消息", "主图初始消息", "子图响应"]
2. 带私有状态的子图
# 主图状态(仅共享字段)
class ParentState(TypedDict):
parent_messages: list
# 子图状态(共享+私有)
class SubgraphState(TypedDict):
parent_messages: list # 共享字段
sub_message: str # 子图私有字段
# 子图节点:修改共享字段+设置私有字段
def subgraph_node(state: SubgraphState) -> SubgraphState:
state["parent_messages"].append("子图修改的共享数据")
state["sub_message"] = "子图私有数据" # 主图无法访问
return state
# 构建子图
sub_builder = StateGraph(SubgraphState)
sub_builder.add_node("sub_node", subgraph_node)
sub_builder.add_edge(START, "sub_node")
sub_builder.add_edge("sub_node", END)
subgraph = sub_builder.compile()
# 主图嵌入子图
parent_builder = StateGraph(ParentState)
parent_builder.add_node("subgraph_node", subgraph)
parent_builder.add_edge(START, "subgraph_node")
parent_builder.add_edge("subgraph_node", END)
parent_graph = parent_builder.compile()
# 执行:主图仅能获取共享字段
result = parent_graph.invoke({"parent_messages": ["主图初始数据"]})
print("主图最终状态:", result) # 仅包含parent_messages,无sub_message
关键注意事项
- 状态合并策略(如 add)会导致主图 + 子图的状态双重拼接,需注意数据重复;
- 子图私有状态仅在子图内部有效,主图无法访问;
- 多团队协作时,子图可由独立团队开发,主图负责整合,降低耦合。
总结
LangGraph 的四大高级特性从「实时性」「可靠性」「灵活性」「复用性」四个维度支撑了复杂工作流的构建:
- 流式处理:让工作流状态实时可见,提升交互体验;
- 状态持久化:保障工作流状态不丢失,支持断点续跑;
- 时间回溯:允许修改历史状态重新执行,适配复杂业务的试错需求;
- 子图:实现工作流模块化复用,降低大型项目的维护成本。
这些特性的组合使用(如「子图 + 持久化 + 流式」)可满足绝大多数生产级多智能体工作流的需求,是 LangGraph 区别于其他工作流框架的核心竞争力。在实际开发中,建议根据业务场景选择合适的特性组合:测试环境用内存持久化 + 基础流式,生产环境用 PostgreSQL 持久化 + 自定义流式 + 子图模块化。




