欢迎光临
我们一直在努力

LangGraph 核心高级特性全解析:流式处理(Streaming)、状态持久化(Persistence)、时间回溯(Time-Travel)与子图(Subgraphs)

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):工作流的「时光机」

核心价值

时间回溯是状态持久化的高阶应用,支持:

  • 回溯到工作流任意历史检查点;
  • 修改历史状态后重新执行(如修正错误的节点输出);
  • 对比不同执行路径的结果,实现「分支实验」。

核心流程

  • 执行工作流并生成检查点;
  • 通过 get_state_history 获取历史检查点列表;
  • (可选)通过 update_state 修改历史状态;
  • 从指定检查点恢复执行,生成新的执行路径。
  • 核心代码示例

    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):工作流的模块化复用

    核心价值

    子图允许将一组节点封装为独立的「子工作流」,并作为单个节点嵌入到主图中,实现:

    • 模块化开发(如多智能体拆分);
    • 工作流复用(相同逻辑无需重复编写);
    • 复杂工作流的分层管理(主图管流程,子图管细节)。

    核心规则

  • 子图与主图的状态需兼容(共享字段名,避免 KeyError);
  • 子图执行时会触发状态合并(默认 add 策略,增量拼接);
  • 子图私有状态不会传递到主图(主图仅能访问共享字段)。
  • 核心代码示例

    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 持久化 + 自定义流式 + 子图模块化。

    赞(0)
    未经允许不得转载:171主机测评 » LangGraph 核心高级特性全解析:流式处理(Streaming)、状态持久化(Persistence)、时间回溯(Time-Travel)与子图(Subgraphs)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址