跳转至

langchain-ai/langgraph:把 Agent 执行流程画成一张图的框架

LangGraph 是 LangChain 团队专门为复杂 Agent 场景做的框架。和 LangChain 那种"从头到尾一条链"的线性写法不同,LangGraph 把整个执行流程建模成一张图:每个节点是一步处理逻辑,节点之间用边声明流转关系,图里可以有回路(一个节点执行完又绕回去)、可以并行、还可以在中途暂停等人工审批。当智能体应用涉及复杂的动态回环、并发协作及长周期挂起(如人工审批)时,传统的线性链(Chain)架构就不够用了——这也是 LangGraph 放弃线性流、改用图结构的原因。搞懂它不能只看 API 文档里怎么连线,得往下挖三层:一轮"超级步"是怎么推进的、多个节点同时写状态时靠什么机制收敛、以及为什么它敢说自己支持"时间旅行"式的回滚。


1. 最小例子:两个节点的图

在讲底层机制之前,先看 LangGraph 最基本的骨架长什么样——不涉及并发、不涉及条件路由,就是两个节点顺序执行一次:

from typing import TypedDict
from langgraph.graph import StateGraph, START, END

class State(TypedDict):
    count: int

def add_one(state: State) -> dict:
    return {"count": state["count"] + 1}

# 1. 定义一份共享状态
graph = StateGraph(State)

# 2. 注册节点:节点就是一个"接收状态、返回状态更新"的函数
graph.add_node("add_one", add_one)

# 3. 用边声明节点之间怎么流转
graph.add_edge(START, "add_one")
graph.add_edge("add_one", END)

# 4. 编译成可以 invoke 的图
app = graph.compile()
print(app.invoke({"count": 0}))  # {'count': 1}

这就是 LangGraph 的基本骨架:定义共享状态(State)、定义处理状态的节点(Node)、用边(Edge)声明流转关系,最后编译(compile)成一张可以调用的图。真正让 LangGraph 区别于普通流程图引擎的,是它底层驱动这张图运行的方式——接下来要讲的 Pregel 超级步模型。


2. 核心设计思路:借鉴 Google Pregel 的"超级步"模型

LangGraph 不是把节点一个个顺序往下跑,而是借鉴了 Google 的 Pregel 分布式图计算框架:所有当前处于活跃状态的节点,在同一轮(称为"超级步",Superstep)里一起并发执行;这一轮跑完之后,系统再根据各节点写入的结果,决定下一轮该激活哪些节点。这样即使多个节点需要并发处理,也有一个清晰的"轮次"边界,方便做状态合并、持久化和后面要讲的回滚。

具体来说,在 LangGraph 中,智能体的运行流程被严格拆解为一系列离散的超级步(Supersteps)

========================================================================================
【Pregel 超级步顺序推进流程】
  超级步 N-1                               超级步 N                               超级步 N+1
 +----------+  (通道消息分发)             +----------+  (通道消息分发)             +----------+
 | Node A/B | ─────────────────────────> |  Node C  | ─────────────────────────> |  Node D  |
 | 并行计算 |  写入 Channel_1/2          | 读取计算 |  写入 Channel_3            | 读取计算 |
 +----------+                            +----------+                            +----------+
      │                                       │                                       │
======V=======================================V=======================================V=========
 【Checkpointer: 每步结束自动执行状态落盘】   【Checkpointer: 自动落盘】                【Checkpointer: 自动落盘】
========================================================================================

2.1 超级步计算原语

  1. 并行节点执行:在每一个超级步 \(N\) 内,所有处于活跃状态的 图节点(Nodes) 共享输入状态并开始并发计算
  2. 通道通信约束:节点之间绝不互相调用或传递变量,而是通过在超级步结束时,向全局 状态通道(Channels) 写入更新数据。
  3. 时序收敛推进:当超级步内所有节点全部执行完毕后,执行引擎统一搜集所有通道的写入数据,利用 Reducer 融合函数 更新图的全局 State,随后进入下一个超级步 \(N+1\)。直到没有节点被激活,或者触发了强制终止条件,图停止运行。

3. 状态怎么合并:Reducer 机制

既然同一个超级步里多个节点可以并发写状态,一个自然的问题是:如果两个节点同时往同一个字段写数据,谁的结果算数?这正是 Reducer 要解决的问题。可以先把它理解成一条"合并规则"——默认情况下,后写入的值会直接覆盖前面的值;但如果业务上需要"追加"而不是"覆盖"(比如多个节点都要往消息列表里塞新消息),就需要显式声明一个 Reducer 函数,告诉框架该怎么合并。

具体机制是:在 LangGraph 中,状态由若干独立通道组成,而不是任意的字典。每一个通道都可以挂载一个 Reducer 函数,定义了当多个节点在同一超级步向该通道写入更新数据时,系统如何合并冲突

from typing import Annotated, Sequence
from typing_extensions import TypedDict

# 使用 Annotated 绑定 Reducer 融合规则
class AgentState(TypedDict):
    # 默认通道:后者写入的值会直接覆盖 (Replace) 前者
    query: str

    # 数组通道:挂载 Reducer,定义合并逻辑为追加 (Append) 而非覆盖
    # 这确保了多节点并行运行时,历史消息不会发生竞态丢失
    messages: Annotated[Sequence[dict], lambda x, y: x + y]
  • 工程价值:通过 Reducer 机制,LangGraph 在机制上解决了并发智能体在同一步骤修改状态时的竞态冲突(Race Conditions),使得大规模并行计算拓扑具备了确定性行为边界。

4. Checkpointer:自动存档与"时间旅行"回滚

超级步模型还带来一个额外的好处:因为整个运行过程被拆成一轮一轮离散的步骤,LangGraph 可以在每一轮结束后自动把当前状态存盘,这样就能像读档一样把状态"倒回"到之前某一步——这也是为什么 LangGraph 敢说自己支持"时间旅行"。这一层机制在 LangGraph 中由 Checkpointer(状态机快照持久化层) 实现,是它最具创新性的工程实践。

4.1 状态持久化与超级步存盘

在每个超级步执行完毕后,Checkpointer 会对当前的图全局状态(各通道的值)执行增量深拷贝与数据库存盘

Database Status:
  - Thread_101, Superstep 0: {"query": "Hello", "messages": [...]}
  - Thread_101, Superstep 1: {"query": "Hello", "messages": [..., {"tool_call": ...}]}
  - Thread_101, Superstep 2: {"query": "Hello", "messages": [..., {"tool_call": ...}, {"tool_result": ...}]}

4.2 时间旅行 (Time-Travel / Rollback) 与分叉运行

由于每个超级步的历史快照都被永久保存,LangGraph 默认支持两大高阶企业模式: * 回滚与恢复(Rollback):如果模型在第 3 步执行了错误的工具导致奔溃,你可以将系统状态回拨到第 2 步的快照,更正输入后继续推进,这被称为时间旅行。 * 交互式审批(Human-in-the-Loop):可以在高危节点前设置 interrupt_before(前置拦截)。此时图执行引擎会在超级步结束前自动将当前 Thread 挂起并持久化。直到人工通过 API 传入确认或修改指令,图才从该快照节点读取现场并"满血恢复"继续向后运行。


5. 进阶示例:一个带条件路由和持久化记忆的 Agent 骨架

有了 Reducer 和 Checkpointer 这两个机制打底,就可以把最开始的两节点最小例子,扩展成一个真正能用的 Agent:能根据大模型的决策走不同分支(条件路由),并且带持久化记忆(Checkpointer)。以下是使用 Python 编写的 LangGraph 生产级智能体拓扑结构骨架,包含了 Reducer 状态定义、节点调用及 Checkpointer 内存持久化绑定:

import operator
from typing import Annotated, TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver

# 1. 定义状态机数据通道及 Reducer 融合策略
class OverallState(TypedDict):
    query: str
    # messages 字段并行更新时,行为定义为数组合并
    messages: Annotated[list, operator.add]

# 2. 定义业务逻辑节点 (Mock 模拟)
def input_node(state: OverallState) -> dict:
    return {"messages": [{"role": "user", "content": state["query"]}]}

def agent_decision_node(state: OverallState) -> dict:
    # 模拟大模型判定:追加助理消息
    user_query = state["messages"][-1]["content"]
    if "db" in user_query.lower():
        return {"messages": [{"role": "assistant", "content": "Query DB tool", "tool_call": "fetch_db"}]}
    return {"messages": [{"role": "assistant", "content": "Reply directly"}]}

def tool_execution_node(state: OverallState) -> dict:
    # 模拟外部数据库工具调用
    return {"messages": [{"role": "tool", "content": "Data Row: 42"}]}

# 3. 编排状态图拓扑关系
workflow = StateGraph(OverallState)

# 注册节点
workflow.add_node("input_loader", input_node)
workflow.add_node("agent", agent_decision_node)
workflow.add_node("db_tool", tool_execution_node)

# 连接拓扑边
workflow.add_edge(START, "input_loader")
workflow.add_edge("input_loader", "agent")

# 条件路由控制
def should_continue(state: OverallState) -> str:
    last_msg = state["messages"][-1]
    if last_msg.get("tool_call") == "fetch_db":
        return "go_to_tool"
    return "go_to_end"

workflow.add_conditional_edges(
    "agent",
    should_continue,
    {
        "go_to_tool": "db_tool",
        "go_to_end": END
    }
)

workflow.add_edge("db_tool", "agent") # 工具结果流回大模型决策节点进行下一轮迭代

# 4. 绑定内存持久化引擎并编译图状态机
memory = MemorySaver()
app = workflow.compile(checkpointer=memory)

# 5. 执行测试
if __name__ == "__main__":
    config = {"configurable": {"thread_id": "session_uuid_10086"}}
    initial_input = {"query": "Find data in DB"}

    # 模拟超级步驱动
    for event in app.stream(initial_input, config):
        for node, state_update in event.items():
            print(f"--- Node '{node}' finished superstep with updates: ---")
            print(state_update)

6. 常见故障与排查

故障模式 底层诱因 系统级现象 预防与排查手段
无限超级步循环 (Infinite Superstep Loop) 状态未收束且没有条件终止判定,导致大模型和工具节点进入无限互拨死循环。 耗尽 Token 额度,API 请求无限挂起,单次请求响应超时报错。 1. 在 compile() 编译时配置 recursion_limit 超级步步数硬限制(默认值 25)。
2. 严格对模型的决策逻辑设置防御条件分支。
Reducer 冲突异常 (Merge Conflict) 并行节点返回的更新数据结构与 Annotated 绑定的 Reducer 输入类型参数不兼容。 系统运行到当前超级步合并阶段抛出 TypeError: ...,进程阻断异常。 规范所有 Node 的返回值,确保它们输出的字典字段类型与 TypedDict 中通道绑定的 Reducer 严格对齐。
持久化锁争抢 (Checkpoint Write Lock) 分布式环境下,针对同一个 thread_id 高并发并发提交任务,导致持久化数据库发生写锁死锁。 客户端收到 Database Locked 错误,或数据库 CPU 飙升。 1. 在 API 网关层实施严格的 Session 级别序列化限流。
2. 对分布式 Checkpointer(如 PostgresCheckpointer)启用带重试的乐观锁或分布式排他锁。

7. 资深系统架构师面试表达方案

面试提问:LangGraph 在构建复杂 Agent 时的底层架构设计是怎样的?它在状态维护和容灾处理上相比于其他工具有什么独特优势?

回答模版: LangGraph 最初吸引我的地方,是它把 Agent 的执行流程建模成 Google Pregel 那套超级步(Superstep)模型,而不是简单的线性 Chain。第一次接触时最难适应的反而是 Reducer 这个概念——节点之间不能互相传值,只能往通道(Channel)里写,写完之后由 Reducer 决定怎么合并。一开始我们图省事,好几个节点都往同一个 messages 通道写,又没绑定 Reducer,结果并发写入时后写的直接覆盖了先写的,历史消息莫名其妙丢了一段,查了半天才想起来要用 Annotated[list, operator.add] 显式声明"追加而不是覆盖"。

真正让我觉得这套设计值回票价的,是 Checkpointer 带来的时间旅行能力。有一次模型在某一步选错了工具、把状态污染了,因为每个超级步的快照都落了盘,我们直接把状态回滚到出错前一步,改完参数继续跑,不用把整个会话重新走一遍。人工审批这类需要中途挂起等确认的场景,也是靠这套快照机制实现的——本质上就是把"执行到哪一步"这件事变成了可以随意读写的数据,而不是一次性跑到底的黑盒。