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 的脑子里外化成了一张图:
这张图有几个关键特性,是它区别于传统 Agent 框架的根本:
- 控制流是显式的:分支条件写在代码里(一个普通函数),不是靠 LLM "猜"。LLM 只负责它擅长的语义理解,流程编排交给确定性代码。
- 状态是共享且可合并的:所有节点读写同一个 State 对象,通过 reducer 决定如何合并(比如消息是追加还是覆盖)。
- 执行是可中断、可恢复的:配合 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),每个超步里:
- 所有被"激活"的节点并行执行;
- 每个节点读取当前 State,计算,产出一个更新;
- 所有更新在超步结束时统一通过 reducer 合并进 State;
- 引擎根据合并后的 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。当一个超步结束,引擎会:
- 找到所有从当前节点出发的边;
- 对条件边,调用路由函数,拿到返回的节点名;
- 把这些节点名加入下一超步的 tasks。
这里有个常被忽略的细节:路由函数必须是纯函数式的(只依赖 State,不产生副作用),因为它在引擎的调度逻辑里被调用,可能被执行多次(比如配合 checkpointer 重放时)。很多人把 print 或者 API 调用写进路由函数,导致重放时行为不一致——这是个大坑,后面会细讲。
4. Checkpointer:状态持久化与时间旅行
Checkpointer 是 LangGraph 从"玩具"走向"生产"的分水岭。它的工作是:在每一个超步结束后,把当前 State 快照持久化。每个快照关联一个 thread_id 和一个 checkpoint_id。
有了它,你可以:
- 断点续跑:进程崩了,传同样的
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 式动态扇出。这些会在你处理更复杂的生产流程时派上大用场。