LangGraph状态图驱动Agent的实现原理

引言

去年我在给一家做智能客服的团队做架构评审时,遇到一个非常典型的问题。他们的 Agent 系统最初用 LangChain 的 AgentExecutor 搭建,几十行代码就能跑通一个"查订单 + 退款"的流程,上线初期效果不错。但随着业务复杂度上升,问题接踵而至:

  • 用户说"帮我查下上周那笔订单,如果还没发货就取消"——这需要先查询、再根据结果条件分支、再执行操作,AgentExecutor 的 ReAct 循环里 Agent 经常"忘记"自己查过什么,或者在中途被 LLM 的自由发挥带偏;
  • 客服流程要求"退款前必须人工审核",但 Agent 有时直接跳过了审核节点;
  • 出了问题想复现某次对话的决策路径,发现日志里只有一串松散的 tool call,根本还原不出控制流;
  • 想给流程加一个"超时重试"和"断点续跑",几乎要把整个 Agent 重写。

这些痛点的本质是:当 Agent 从"一问一答"进化到"多步骤、有条件、有状态、需要人工介入"的复杂流程时,靠 LLM 自由编排控制流是不可靠的。你需要的不是更聪明的 Prompt,而是一套显式的、可持久化的、可被程序精确控制的状态机。

LangGraph 就是为解决这个问题而生的。它的核心思想可以用一句话概括:把 Agent 的执行过程建模成一张有向图,节点是计算单元,边是控制流,而整张图共享一份可持久化的状态。这篇文章我会从源码级别拆解它的实现原理,并给出可直接运行的实战代码。


核心概念:用"餐厅后厨"理解状态图

先别急着看代码。想象一家餐厅的后厨:

  • State(状态) 是那张流水线订单票。从"下单"到"出餐",这张票在后厨各个工位之间传递,上面记录着:点了什么菜、做到哪一步了、有没有特殊要求。每个工位都能读到它、修改它。
  • Node(节点) 是每个工位:切菜台、炒锅、摆盘台。每个工位是一个独立的"处理单元",接收订单票,干活,更新票上的信息。
  • Edge(边) 是工位之间的传送带:切完菜传给炒锅,这是普通边;但如果订单里是"凉菜",就绕过炒锅直接去摆盘台,这就是条件边。
  • Checkpointer(检查点) 是后厨的监控录像 + 存档。每经过一个工位,系统自动快照一次订单票的状态。厨师突然晕倒了?换个厨师,从上次快照继续做。顾客中途要加菜?从当前状态改一下继续。

对比一下 LangChain 的 AgentExecutor:那更像是一个厨师包办所有事。他既切菜又炒菜又摆盘,你觉得他累了想换人,或者想在中途插一个"质检员",就很难办——因为流程在他脑子里,不在系统里。

而 LangGraph 把"流程"从 LLM 的脑子里外化成了一张图:

graph TD A[__start__] --> B[理解意图节点] B --> C{路由判断} C -->|查询类| D[查询订单节点] C -->|操作类| E[执行操作节点] D --> F{是否可退款} F -->|是| G[人工审核节点] F -->|否| H[返回结果节点] G --> H E --> H H --> I[__end__]

这张图有几个关键特性,是它区别于传统 Agent 框架的根本:

  1. 控制流是显式的:分支条件写在代码里(一个普通函数),不是靠 LLM "猜"。LLM 只负责它擅长的语义理解,流程编排交给确定性代码。
  2. 状态是共享且可合并的:所有节点读写同一个 State 对象,通过 reducer 决定如何合并(比如消息是追加还是覆盖)。
  3. 执行是可中断、可恢复的:配合 Checkpointer,图可以在任意节点暂停,持久化,之后从断点继续——这是"人工介入(Human-in-the-loop)"的实现基础。

源码级原理深度分析

这一节是重点。我会基于 LangGraph 的核心源码(以 langgraph 的 Python 实现为主,Pregel 执行模型)来拆解它是怎么跑起来的。

1. 状态契约:TypedDict + Annotated + Reducer

LangGraph 的 State 本质上是一个 TypedDict,而每个字段可以带一个 Annotated 注解来指定 reducer(归约函数),决定多个节点更新同一个字段时如何合并。

from typing import TypedDict, Annotated
from operator import add
from langgraph.graph.message import add_messages

class AgentState(TypedDict):
    # 消息列表:用 add_messages 作为 reducer
    # 语义是"追加",而不是"覆盖"
    messages: Annotated[list, add_messages]
    # 计数器:用 operator.add 作为 reducer
    step_count: Annotated[int, add]
    # 普通字段:没有 Annotated,默认行为是"后写覆盖"
    current_intent: str

为什么需要 reducer?这是 LangGraph 状态模型最关键的设计之一。在并发或多次更新场景下,如果每个节点都返回完整的新 State,会带来两个问题:一是节点必须知道"别人改了什么",耦合太强;二是无法表达"追加""累加"这类语义。

源码视角:在 LangGraph 内部,State 被抽象成 StateGraph 持有一个 channels 字典。每个字段对应一个 Channel(通道),Channel 决定了如何接收和合并更新:

  • 没有 reducer 的字段 → 对应 LastValue 通道,新值直接覆盖旧值;
  • Annotated[list, add_messages] → 对应 Topic/BinaryOperatorAggregate 通道,新值通过 reducer 函数与旧值合并。

当节点函数返回 {"messages": [new_msg]} 这样一个部分更新(partial update) 时,执行引擎并不会把它当成完整 State,而是把它喂给对应 Channel 的 reducer,与当前 State 合并。这就是为什么节点可以只返回它关心的那部分字段。

2. Pregel 执行模型:超步(Super-step)与消息传递

LangGraph 的执行引擎借鉴了 Google 的 Pregel 图计算模型。理解这一点,你就理解了它的调度本质。

Pregel 的核心是 BSP(Bulk Synchronous Parallel,整体同步并行) 模型。执行被切成一个个 超步(Super-step),每个超步里:

  1. 所有被"激活"的节点并行执行;
  2. 每个节点读取当前 State,计算,产出一个更新;
  3. 所有更新在超步结束时统一通过 reducer 合并进 State;
  4. 引擎根据合并后的 State 和边定义,决定下一个超步激活哪些节点。

用生活类比:这就像快递分拣中心的传送带。每一"轮"传送带上,所有该处理的包裹被并行处理;处理完的包裹统一汇入下一段传送带;系统根据每个包裹的目的地(条件边)决定它流向哪个分拣口。关键是同一轮里大家看到的都是本轮开始时的状态,不会出现"我读的时候你还没写完"的竞态。

源码视角:核心调度逻辑在 langgraph/pregel/__init__.py 的 PregelLoop 中。它维护一个 loop.tasks 集合,代表当前超步待执行的节点。每一轮循环大致是:

while loop.tasks:                      # 还有待执行任务
    # 1. 并行执行所有 task,收集 writes(写入)
    # 2. 将 writes 通过 channels 的 reducer 合并进 state
    # 3. 根据 state 和边,计算出下一轮的 tasks

这里有个精妙之处:节点产出的不是直接修改全局 State,而是写到一个 writes 缓冲里,超步结束时才 apply。这保证了超步内的确定性——并行节点之间不会互相看到中间状态,避免了难以调试的时序问题。

3. 条件边与路由:控制流如何被"编译"

条件边(add_conditional_edges)的机制是:你提供一个路由函数,它接收当前 State,返回一个字符串(或多个字符串),引擎据此决定激活哪个后继节点。

源码视角:路由函数在编译期被包装成一个特殊的 task。当一个超步结束,引擎会:

  1. 找到所有从当前节点出发的边;
  2. 对条件边,调用路由函数,拿到返回的节点名;
  3. 把这些节点名加入下一超步的 tasks。

这里有个常被忽略的细节:路由函数必须是纯函数式的(只依赖 State,不产生副作用),因为它在引擎的调度逻辑里被调用,可能被执行多次(比如配合 checkpointer 重放时)。很多人把 print 或者 API 调用写进路由函数,导致重放时行为不一致——这是个大坑,后面会细讲。

4. Checkpointer:状态持久化与时间旅行

Checkpointer 是 LangGraph 从"玩具"走向"生产"的分水岭。它的工作是:在每一个超步结束后,把当前 State 快照持久化。每个快照关联一个 thread_id 和一个 checkpoint_id。

graph LR S0[初始状态] -->|超步1| S1[快照1] S1 -->|超步2| S2[快照2] S2 -->|超步3| S3[快照3] S3 -.中断.-> H[人工介入] H -->|恢复| S3 S3 -->|超步4| S4[快照4]

有了它,你可以:

  • 断点续跑:进程崩了,传同样的 thread_id,从最后一个快照继续;
  • 时间旅行:传一个历史 checkpoint_id,从任意历史状态重新分叉执行;
  • 人工介入:在某个节点前 interrupt,等人审批后 resume;
  • 多轮对话:同一个 thread_id 天然就是一段会话记忆。

源码视角:Checkpointer 是一个接口(BaseCheckpointSaver),内置实现有 MemorySaver(内存)、SqliteSaver、PostgresSaver。核心方法有四个:put(写快照)、get_tuple(读快照)、list(列历史)、put_writes(写中间写入)。生产环境一定要用 PostgresSaver 这类外部存储,MemorySaver 重启即丢。

5. 状态更新合并的时序细节

再深挖一层。假设一个超步里,节点 A 和节点 B 都被激活,它们都更新了 messages 字段。会发生什么?

由于 messages 用的是 add_messages reducer,合并逻辑是按顺序追加。LangGraph 内部给每个写入分配了 (task_id, idx) 这种可排序的标识,保证合并顺序是确定性的——即使 A、B 并行执行,它们的写入合并顺序也是可复现的。这个设计直接支撑了"时间旅行重放"的正确性:同样的输入,重放会得到同样的结果。

这点非常关键。如果你自己实现一个并行 Agent 调度器,用 asyncio.gather 收集结果然后 list.extend,你会发现顺序是不确定的(取决于谁先完成),导致同样的对话重放两次结果不一样。LangGraph 通过给写入排序解决了这个经典难题。


实战代码

下面给三个完整可运行的示例,层层递进。

示例 1:最小可运行的状态图 Agent

"""
最小状态图 Agent:理解意图 -> 条件路由 -> 执行 -> 结束
依赖: pip install langgraph langchain-openai
"""
from typing import TypedDict, Literal
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from typing import Annotated
from langchain_openai import ChatOpenAI

# ---------- 1. 定义状态契约 ----------
class State(TypedDict):
    # add_messages reducer:保证消息是"追加"语义,而非覆盖
    messages: Annotated[list, add_messages]
    # 意图识别结果,普通字段,后写覆盖
    intent: str

# ---------- 2. 初始化 LLM ----------
llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

# ---------- 3. 定义节点 ----------
def classify_intent(state: State) -> dict:
    """节点:让 LLM 判断用户意图,返回部分更新"""
    user_msg = state["messages"][-1].content
    # 用一个极简 prompt 做分类,真实场景可以用结构化输出
    resp = llm.invoke(
        f"判断下面这句话的意图,只回复'query'或'action'其中一个词:\n{user_msg}"
    )
    intent = resp.content.strip().lower()
    # 只返回需要更新的字段,引擎会通过 reducer 合并
    return {"intent": intent}

def handle_query(state: State) -> dict:
    """节点:处理查询类意图"""
    return {"messages": [("assistant", "已为您查询到订单信息:状态为已发货。")]}

def handle_action(state: State) -> dict:
    """节点:处理操作类意图"""
    return {"messages": [("assistant", "操作已受理,正在为您取消订单。")]}

# ---------- 4. 定义路由函数(必须是纯函数) ----------
def route_by_intent(state: State) -> Literal["handle_query", "handle_action"]:
    """根据意图决定走哪条分支"""
    if state["intent"] == "query":
        return "handle_query"
    return "handle_action"

# ---------- 5. 组装图 ----------
builder = StateGraph(State)
builder.add_node("classify_intent", classify_intent)
builder.add_node("handle_query", handle_query)
builder.add_node("handle_action", handle_action)

builder.add_edge(START, "classify_intent")
# 条件边:从 classify_intent 出发,由 route_by_intent 决定去向
builder.add_conditional_edges(
    "classify_intent",
    route_by_intent,
    # 显式声明可能的目标,便于引擎校验和可视化
    {"handle_query": "handle_query", "handle_action": "handle_action"},
)
builder.add_edge("handle_query", END)
builder.add_edge("handle_action", END)

# 编译成可执行图
graph = builder.compile()

# ---------- 6. 运行 ----------
if __name__ == "__main__":
    result = graph.invoke(
        {"messages": [("user", "帮我查一下我的订单")]}
    )
    print(result["messages"][-1].content)
    # 输出: 已为您查询到订单信息:状态为已发货。

这个例子展示了骨架:State 契约、节点、条件边、编译运行。注意节点返回的都是部分更新,而不是完整 State。

示例 2:带 Checkpointer 的多轮对话 + 人工介入

"""
带持久化和人工介入的 Agent
依赖: pip install langgraph langchain-openai
"""
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command
from langchain_openai import ChatOpenAI

class State(TypedDict):
    messages: Annotated[list, add_messages]
    # 记录审批结果
    approved: bool

llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

def draft_refund(state: State) -> dict:
    """节点:起草退款方案"""
    return {"messages": [("assistant", "拟退款金额 ¥299,请审批。")]}

def human_approval(state: State) -> dict:
    """
    节点:人工审批。
    interrupt 会在此处暂停图执行,把控制权交还给调用方。
    恢复时,resume 传入的值会作为 interrupt() 的返回值。
    """
    decision = interrupt("是否批准本次退款?(yes/no)")
    approved = str(decision).lower() == "yes"
    return {"approved": approved}

def execute_refund(state: State) -> dict:
    if state.get("approved"):
        return {"messages": [("assistant", "退款已执行。")]}
    return {"messages": [("assistant", "退款被驳回。")]}

builder = StateGraph(State)
builder.add_node("draft_refund", draft_refund)
builder.add_node("human_approval", human_approval)
builder.add_node("execute_refund", execute_refund)

builder.add_edge(START, "draft_refund")
builder.add_edge("draft_refund", "human_approval")
builder.add_edge("human_approval", "execute_refund")
builder.add_edge("execute_refund", END)

# 关键:挂上 checkpointer,interrupt/resume 才能工作
memory = MemorySaver()
graph = builder.compile(checkpointer=memory)

if __name__ == "__main__":
    config = {"configurable": {"thread_id": "user-42-order-1001"}}

    # 第一次运行:会在 human_approval 处中断
    graph.invoke({"messages": [("user", "我要退款")]}, config)
    snapshot = graph.get_state(config)
    print("中断点:", snapshot.next)  # ('human_approval',)

    # 恢复执行,传入审批结果
    graph.invoke(Command(resume="yes"), config)
    final = graph.get_state(config)
    print("最终输出:", final.values["messages"][-1].content)
    # 输出: 退款已执行。

这个例子的核心是 interrupt + Command(resume=...)。注意 thread_id 就是这段会话的"存档位",同一个 thread_id 可以在任意时刻恢复。生产环境把 MemorySaver 换成 PostgresSaver 即可。

示例 3:并行节点 + 自定义 Reducer(带错误聚合)

"""
并行执行多个检查节点,并用自定义 reducer 聚合结果
依赖: pip install langgraph
"""
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END

# ---------- 自定义 reducer:合并多个检查结果 ----------
def merge_checks(left: dict, right: dict) -> dict:
    """
    把多个并行节点的检查结果合并成一个字典。
    注意:reducer 必须是纯函数、可交换且幂等的(在合并语义上)。
    """
    return {**left, **right}

class State(TypedDict):
    user_input: str
    # 用自定义 reducer,把并行节点的 dict 结果合并
    checks: Annotated[dict, merge_checks]

# ---------- 三个可并行执行的检查节点 ----------
def check_profanity(state: State) -> dict:
    ok = "脏话" not in state["user_input"]
    return {"checks": {"profanity": ok}}

def check_length(state: State) -> dict:
    ok = len(state["user_input"]) <= 200
    return {"checks": {"length": ok}}

def check_spam(state: State) -> dict:
    ok = "http://" not in state["user_input"]
    return {"checks": {"spam": ok}}

def summarize(state: State) -> dict:
    """汇总节点:读取合并后的 checks"""
    failed = [k for k, v in state["checks"].items() if not v]
    verdict = "通过" if not failed else f"未通过: {failed}"
    print(f"审核结论: {verdict}")
    return {}

builder = StateGraph(State)
builder.add_node("check_profanity", check_profanity)
builder.add_node("check_length", check_length)
builder.add_node("check_spam", check_spam)
builder.add_node("summarize", summarize)

# 从 START 扇出到三个检查节点 —— 它们会在同一个超步并行执行
builder.add_edge(START, "check_profanity")
builder.add_edge(START, "check_length")
builder.add_edge(START, "check_spam")

# 三个检查节点扇入到 summarize
builder.add_edge("check_profanity", "summarize")
builder.add_edge("check_length", "summarize")
builder.add_edge("check_spam", "summarize")
builder.add_edge("summarize", END)

graph = builder.compile()

if __name__ == "__main__":
    graph.invoke({"user_input": "这是一条正常留言"})
    # 输出: 审核结论: 通过
    graph.invoke({"user_input": "这是脏话内容 http://spam.com"})
    # 输出: 审核结论: 未通过: ['profanity', 'spam']

这个例子演示了 扇出/扇入(fan-out/fan-in) 模式。从 START 出发的三条边让三个检查节点在同一超步并行执行,引擎会等它们全部完成后,再用 merge_checks 把结果合并,然后统一触发 summarize。这正是 Pregel BSP 模型的直观体现。


方案对比:LangGraph vs 其他 Agent 编排方案

方案 控制流 状态管理 持久化/中断 适用场景
LangChain AgentExecutor LLM 自由 ReAct 循环 松散,靠 memory 弱 简单问答、单轮工具调用
LangGraph 显式状态图 共享 State + Reducer 强(checkpointer) 复杂多步、需人工介入、需可复现
AutoGen 多 Agent 对话驱动 会话消息 中 多智能体协作、群聊式
CrewAI 角色/任务编排 任务级 中 角色分工明确的流水线
手写状态机 完全自定义 完全自定义 自己实现 流程固定、极致可控

几个关键判断:

  • 如果你的流程是"一次工具调用就结束",别用 LangGraph,AgentExecutor 或直接 function calling 足够,引入状态图是过度设计。
  • 如果流程需要"条件分支 + 人工审批 + 断点续跑 + 可复现",LangGraph 几乎是目前 Python 生态里最成熟的选择。
  • 如果是"多个平等 Agent 互相讨论",AutoGen 的对话模型更自然;LangGraph 也能做(用 messages 状态 + 循环边),但要多写一些编排代码。
  • 如果流程极其固定、对延迟极度敏感,手写状态机反而更轻量,LangGraph 的抽象层会带来一点开销。

LangGraph 的独特价值在于:它把"Agent 流程"当成一等公民来建模,而不是当成 LLM 的副产品。这是它和"更聪明的 Prompt"路线的根本分野。


最佳实践与避坑指南

1. Reducer 选错是头号坑。 消息列表务必用 add_messages,否则每次节点返回都会把历史消息覆盖掉,多轮对话直接失忆。反过来,如果你确实想要"覆盖"语义,就别加 Annotated。

2. 路由函数必须是纯函数。 不要在条件边路由函数里调用 API、写数据库、打日志依赖外部状态。因为它可能在重放时被多次调用,副作用会导致行为不一致。所有副作用放到节点里。

3. 节点返回"部分更新",不要返回完整 State。 返回完整 State 会绕过 reducer 的语义,还容易覆盖掉其他并行节点的写入。养成"只返回我改动的字段"的习惯。

4. 并行节点的写入顺序是确定性的,但别依赖它做业务逻辑。 引擎会给写入排序保证可复现,但你的业务不应该假设"A 一定在 B 之前合并"。需要顺序就串行连边。

5. 生产环境必须用外部 Checkpointer。 MemorySaver 只适合本地调试。PostgresSaver 是生产首选,注意定期清理老 checkpoint,否则表会膨胀。

6. thread_id 的设计就是会话隔离的设计。 一个用户一次会话一个 thread_id,别复用。跨会话共享状态要用其他机制(如外部 KV),不要靠 thread_id。

7. 控制图规模。 节点太多、边太乱时,可视化(graph.get_graph().draw_mermaid())会帮你发现"绕圈"和"死路"。给循环边一定要配终止条件,否则 Agent 会在两个节点间无限循环烧 token。

8. 用 stream_mode 做可观测性。 graph.stream(..., stream_mode="updates") 能实时吐出每个节点的更新,接上日志系统就是天然的决策链路追踪,比事后翻 tool call 日志强太多。


总结

回到开头那个客服团队的问题。他们最终把流程迁到了 LangGraph:意图识别、查询、条件分支、人工审核、执行,全部显式建模成节点和边;用 Postgres Checkpointer 做持久化,退款这种敏感操作走 interrupt 人工审批;出了问题直接用 thread_id + checkpoint_id 重放,决策路径一目了然。代码量没有暴增,但可控性和可维护性完全是两个量级。

回顾本文的核心要点:

  • LangGraph 的本质是 Pregel 风格的状态图执行引擎,把 Agent 流程从 LLM 的"脑子"里外化成显式的图;
  • State + Reducer 是共享状态的契约,决定了多个节点更新如何合并,add_messages 是最常用的 reducer;
  • 超步(BSP)模型 保证了并行执行的确定性和可复现性,这是"时间旅行"和"断点续跑"的基石;
  • Checkpointer 是生产化的分水岭,支撑人工介入、断点续跑、多轮记忆;
  • 路由函数要纯、节点返回部分更新、生产用外部 Checkpointer——这几条是避坑重点。

延伸思考:状态图这套抽象其实不限于 Agent。任何"多步骤、有状态、需要人工介入、要求可复现"的业务流程(工作流引擎、审批流、数据管道)都能用同样的模型表达。LangGraph 真正的贡献,是让 LLM 的语义能力与确定性流程编排各司其职——LLM 负责它擅长的理解与生成,状态图负责它不擅长的流程控制。想清楚这条边界,你就能设计出既聪明又可靠的 Agent 系统。

下一步值得研究的方向:子图(subgraph)组合、多 Agent 之间的状态传递、以及 Send API 实现的 Map-Reduce 式动态扇出。这些会在你处理更复杂的生产流程时派上大用场。