LangGraph 源码地图与运行机制
LangGraph 源码地图与运行机制
这篇文章记录了我对 LangGraph 的一次系统性源码阅读:从用户写下 StateGraph、add_node、compile() 开始,一路追到 invoke() 背后的 Pregel loop、任务调度、LLM 调用和 checkpoint 持久化。阅读范围以 Python libs/langgraph 核心运行时和 libs/prebuilt 的经典 ReAct Agent 为主,并覆盖 checkpoint、stream、store 的边界。
阅读方法:先用符号调用图建立整体结构,再对关键函数按源码顺序逐段阅读。
1. 先给结论:LangGraph 到底是什么
LangGraph 是一个以状态、通道(channel)、任务、超步(superstep)和检查点为核心的图运行时。用户写的是节点函数和边,运行时真正执行的是一组 PregelExecutableTask;每个任务读取当前 checkpoint 的通道快照,产出对通道的 writes,所有任务完成后统一应用 writes,再决定下一轮任务。
最重要的架构分层如下:
flowchart TB
APP[用户代码\nStateGraph / Functional API] --> BUILD[图构建层\nStateGraph]
BUILD --> COMPILED[编译产物\nCompiledStateGraph]
COMPILED --> PREGEL[执行接口\nPregel.invoke / stream]
PREGEL --> LOOP[运行控制\nSyncPregelLoop / AsyncPregelLoop]
LOOP --> ALGO[算法层\nprepare_next_tasks / apply_writes]
LOOP --> RUNNER[任务执行\nPregelRunner / retry]
RUNNER --> NODE[PregelNode.bound\nRunnable / 用户节点 / ToolNode]
NODE --> LLM[LangChain ChatModel\ninvoke / ainvoke]
LOOP --> CP[Checkpoint Saver]
LOOP --> EVENTS[stream handlers\nvalues / updates / messages / debug]
PRE[prebuilt create_react_agent] --> BUILD
PRE --> TOOL[ToolNode]
TOOL --> LOOP
由此可以得到四个关键判断:
StateGraph只负责声明和编译,不负责执行用户节点。invoke()不是独立的执行器;它消费stream()产生的事件并收集最终结果。- LLM 不是 LangGraph 内部的特殊对象。对运行时来说,LLM 只是一个 Runnable,和普通节点一样被包装成任务;真正的模型 HTTP 请求发生在 LangChain provider 实现中。
- Agent 的”思考-调用工具-再思考”不是用 Python
while写成的业务循环,而是图边让agent和tools在连续的 Pregel step 中反复成为下一轮任务。
2. 仓库源码地图
| 层 | 目录 | 主要职责 | 首要入口 |
|---|---|---|---|
| 图 DSL | libs/langgraph/langgraph/graph/ |
schema、节点、边、条件分支、编译 | graph/state.py |
| 核心执行 | libs/langgraph/langgraph/pregel/ |
Runnable 接口、loop、任务、并发、重试、输出 | pregel/main.py |
| 状态算法 | libs/langgraph/langgraph/pregel/_algo.py |
任务准备、触发器匹配、writes 应用 | prepare_next_tasks、apply_writes |
| 任务执行 | libs/langgraph/langgraph/pregel/_runner.py |
并发提交、等待、commit、错误处理 | PregelRunner.tick |
| 运行上下文 | libs/langgraph/langgraph/runtime.py |
context、store、stream_writer、control |
Runtime |
| Runnable 包装 | libs/langgraph/langgraph/_internal/_runnable.py |
函数、Runnable、序列组合与 callback | RunnableCallable、RunnableSeq |
| 预构建 Agent | libs/prebuilt/langgraph/prebuilt/ |
ToolNode、ReAct 图工厂 | chat_agent_executor.py |
| checkpoint 契约 | libs/checkpoint/langgraph/checkpoint/ |
checkpoint 数据结构和 saver 抽象 | checkpoint/base/__init__.py |
| SQLite 后端 | libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/ |
checkpoint 和 pending writes 的 SQLite 实现 | SqliteSaver、AsyncSqliteSaver |
| PostgreSQL 后端 | libs/checkpoint-postgres/langgraph/checkpoint/postgres/ |
服务化、多连接 checkpoint | PostgresSaver |
依赖方向是 checkpoint -> langgraph/prebuilt -> 应用图;SQLite/Postgres 只是 checkpoint 的替换实现。CLI 负责发现和装载用户导出的图,SDK 负责远程 HTTP/stream,不会替代进程内的 Pregel loop。
3. 从用户代码到编译产物
3.1 StateGraph.__init__:先把类型变成通道
文件:libs/langgraph/langgraph/graph/state.py。
初始化按源码顺序做这些事:
- 处理已废弃的
config_schema、input、output参数,并把它们转换到context_schema、input_schema、output_schema。 - 初始化
nodes、edges、branches、schemas、channels、managed、waiting_edges。 - 记录状态、输入、输出、运行上下文三种 schema。没有单独输入/输出 schema 时,它们回退到状态 schema。
- 对 state schema 调用
_add_schema();对 input/output schema 调用同一函数但禁止 managed channel。
_add_schema() 不是简单保存类型:
_get_channels(schema)从 TypedDict、dataclass、Pydantic 或注解中解析普通 channel 和 managed value。- 普通 channel 放入
self.channels,它描述字段的值类型和聚合规则。 - managed value 放入
self.managed,它不是每一步普通状态更新的一部分,例如步数管理器或外部资源。 - 同名 channel 如果类型不同会报错;
LastValue有兼容处理。 - 因此,状态字典只是用户看到的形状,运行时内部真正维护的是”channel 名 -> channel 实例”。
一个字段的更新路径是:
1 | |
3.2 add_node():保存节点规范,不运行节点
具体实现的职责可以按源码执行顺序理解为:
- 识别调用形式:
add_node(fn)、add_node("name", fn)或传入 Runnable。 - 推导或校验节点名,拒绝重复节点名以及与
START/END冲突的名字。 coerce_to_runnable()将普通函数、异步函数或 Runnable 统一成可调用对象。- 保存节点输入 schema、metadata、
defer、retry policy、cache policy、timeout、error handler 和 destinations。 - 如果节点使用了额外 state schema,会再次
_add_schema(),让节点可以只读取状态的一个视图。
这里有一个容易误解的设计:destinations 主要用于渲染/静态说明;真正的运行时路由来自普通边、条件分支、节点返回的 Command 或 Send。
3.3 add_edge() 与 add_conditional_edges():声明触发关系
无条件边被保存为 (source, target)。START、END 是虚拟边界:
1 | |
条件边保存的是 branch:
1 | |
条件函数不是在 add_conditional_edges() 时执行,而是在 source 节点的 writes 已经提交、下一轮准备任务时执行。因此路由看到的是该超步结束后的状态,而不是 source 节点执行的中间态。
3.4 StateGraph.compile():把声明固化成 Pregel
compile() 的核心结果是 CompiledStateGraph,它继承/组合 Pregel 的执行能力。按源码语义,编译过程分为:
- 校验 graph 是否已经编译、输入/输出 channel 是否存在、边的 source/target 是否有效。
- 合并用户传入的 checkpointer、store、cache、interrupt、debug 和 name。
- 创建
CompiledStateGraph,把 builder 的 channels、managed、节点配置、输入输出 channel 交给它。 attach_node()为每个用户节点建立PregelNode:节点的 bound runnable、读取哪些 channel、写入哪些 channel、writers、retry/cache/error 策略都会在这里固定。attach_edge()把普通边编译成 channel trigger。目标节点不是直接被 Python 函数调用,而是被某个 channel 更新触发。attach_branch()把条件分支编译成 branch runnable 和动态写入逻辑。- 设置
START的输入写入和END的终止语义。 - 设置
trigger_to_nodes,这是运行时从”哪些 channel 更新了”反查”哪些节点需要执行”的索引。
所以编译前后是一次重要的表示转换:
1 | |
编译不会调用模型、工具或用户节点;只会包装和连接它们。
4. 一次 invoke() 的完整调用链
主入口:libs/langgraph/langgraph/pregel/main.py。
4.1 invoke() 只是 stream 的收集器
Pregel.invoke() 的实际策略是:
- 根据
stream_mode选择收集方式,默认使用values。 - 调用
self.stream(...)。 - 逐个消费 stream chunk;v1/v2 对 interrupt 的读取方式不同。
values模式返回最新完整状态;其他模式返回最后一个 chunk 或收集后的 chunks。
这就是为什么排查执行问题时应该先看:
1 | |
而不是只盯着 invoke() 的最终字典。
4.2 Pregel.stream() 的逐段流程
阶段 A:解析运行参数
- 处理旧的
checkpoint_during,转换为 durability。 - 没有显式
stream_mode时,顶层使用实例默认值;被另一个图作为节点调用时默认values。 - 创建
SyncQueue,所有 stream handler 最终都向这个队列写入。 ensure_config(self.config, config)合并调用配置和图默认配置。_defaults()解析 stream modes、output keys、interrupt 前后节点、checkpointer、store、cache、durability。- 创建 callback manager 和 graph callback manager,建立一次运行的 run id。
阶段 B:装配流和 Runtime
根据 stream mode 安装 handler:
| mode | handler/来源 | 产生什么 |
|---|---|---|
values |
loop 的 output emitter | 每步完整输出状态 |
updates |
task/output 映射 | 节点名和节点局部 writes |
messages |
StreamMessagesHandler |
LLM token 与 metadata |
tools |
StreamToolCallHandler |
工具调用事件 |
custom |
stream_writer 写入 queue |
节点主动发送的数据 |
tasks/debug/checkpoints |
loop debug emitter | 任务、checkpoint、诊断信息 |
随后创建 Runtime:
1 | |
它被放进 config["configurable"],节点通过参数 runtime 取得,而不是从全局变量读取。
阶段 C:创建 loop 和 runner
SyncPregelLoop(...) 注入:输入、stream protocol、config、store、cache、checkpointer、nodes、channel specs、输入输出 keys、interrupt、retry/cache policy 和 trigger_to_nodes。
然后创建:
1 | |
PregelRunner 不决定图结构;它只负责把已经选好的任务执行起来,并把任务结果 commit 回 loop。
阶段 D:主 BSP loop
源码中的核心循环是:
1 | |
这段代码的每一行对应一个架构保证:
loop.tick():根据当前 checkpoint、pending writes 和更新过的 channel 准备本轮任务;若无任务则结束。match_cached_writes():缓存命中时直接复用任务 writes,不重复执行节点。runner.tick(...):执行本轮所有writes为空的任务;每个 task 可以同步或并发运行。_output(...):把任务完成期间积累的 queue 事件转换成用户可见 chunk。loop.after_tick():统一应用本轮 writes、推进 channel version、写 checkpoint、检查interrupt_after。syncdurability 等待 checkpoint 真正完成;async允许 checkpoint 写入和下一步重叠;exit延迟到退出。
BSP 关键点:第 N 步的 writes 不会被第 N 步其他节点读取,只在 after_tick() 后成为第 N+1 步的 channel 值。 这保证并行节点读取同一个不可变快照,避免执行顺序影响结果。
阶段 E:退出和异常
循环退出后继续吐尽 stream queue,再根据 loop status 处理:
out_of_steps:超过recursion_limit,抛GraphRecursionError。draining:收到RunControl的协作式关闭请求,抛GraphDrained。interrupt_before/interrupt_after:抛GraphInterrupt,checkpoint 已经留下恢复所需信息。- 正常退出:
run_manager.on_chain_end(loop.output)。 - 任意异常:
run_manager.on_chain_error(e)后原样抛出。
5. loop 内部:任务怎样被生成、运行和提交
5.1 SyncPregelLoop.__init__
文件:pregel/_loop.py。
初始化不是”创建一个 while 状态变量”这么简单,它还处理恢复边界:
- 保存输入、nodes、specs、checkpointer、策略和当前 step/stop。
- 判断是否 nested graph、是否 replaying。
- 根据 scratchpad 的 subgraph counter 修正
checkpoint_ns。 - 顶层图清理不应继承的 namespace/id;子图保留自己的命名空间。
- 从
checkpoint_map选择本次运行实际恢复的 checkpoint。 - 把
thread_id规范化为字符串,计算checkpoint_nstuple。 - 从 Runtime 取出
control,供 drain 使用。
这解释了为什么同一个图对象可以安全地被多个 thread 调用:运行状态在 loop 和 checkpoint config 中,不在编译对象的可变业务状态里。
5.2 tick():一轮开始前做什么
- 如果
step > stop,设置out_of_steps并返回False。 - 调用
_algo.prepare_next_tasks(...),传入 checkpoint、pending writes、nodes、channels、managed、step、触发器、retry/cache policy。 - 如果上一轮 checkpoint 写入已经完成,发出 checkpoint debug 事件。
- 没有 tasks 时设置
done并返回False。 - 如果 control 要求 drain,设置
draining并返回False。 - 恢复场景下,把已经成功持久化的 writes 重新挂到对应 task;跳过 ERROR、INTERRUPT、RESUME 等控制写入。
- 检查
interrupt_before。 - 发出 task debug 事件,并输出缓存命中的 writes。
- 返回
True,允许 runner 执行任务。
5.3 prepare_next_tasks():从 channel 更新反查节点
文件:pregel/_algo.py。
该函数的概念算法是:
1 | |
task 中最重要的不是函数本身,而是三组数据:
triggers:为什么该节点本轮应该执行。writes:已完成或从 checkpoint 恢复的结果;为空才交给 runner 执行。bound/node:真正被 Runnable 调用的对象。
5.4 PregelRunner.tick():执行任务并 commit
文件:pregel/_runner.py。
执行顺序如下:
- 把 tasks 转成 tuple,建立
FuturesDict,每个 future 完成时回调commit。 - 单任务、无 timeout、无 waiter 时走 fast path,直接
run_with_retry()。 - 多任务时通过
submit()并发提交,每个任务都携带_call,允许节点内部动态Send新任务。 - 使用
FIRST_COMPLETED等待;一个任务完成就给上层一次 yield 机会,于是 stream 可以边执行边输出。 - 节点错误若有 graph-level error handler,则标记为 handled,并调度 handler task;否则进入 panic/reraise。
- 所有 future 完成后
_panic_or_proceed()决定是否抛出未处理异常。
commit() 的关键语义是:节点执行结果不会直接修改共享 channels,而是转换为 task writes,再调用 loop 的 put_writes(task_id, writes)。因此并发任务之间没有”谁先写谁覆盖”的 Python 字典竞争,冲突由 channel reducer 在超步边界处理。
5.5 after_tick():超步提交点
源码顺序非常关键:
- 收集所有 task writes。
- 调用
apply_writes(...),将本轮 writes 应用到 channels,并返回下一轮需要关注的 updated channels。 - 如果输出 channel 被更新,发出
values。 - 保存 exit durability 需要的 delta writes。
- 清空 pending writes,关闭 replay 状态。
_put_checkpoint({"source": "loop"})保存新的通道版本和 metadata。- 检查
interrupt_after。 - 清掉 resuming 标记。
因此 checkpoint 是”超步提交”的持久化镜像,而不是每个节点函数返回一次就立即替换整个 state。
6. LLM 调用:prompt 如何拼接、模型如何真正被调用
以下以 libs/prebuilt/langgraph/prebuilt/chat_agent_executor.py 的 create_react_agent() 为例。
6.1 Agent 工厂初始化阶段
create_react_agent() 初始化阶段按顺序做:
- 处理 deprecated kwargs,校验
version只能是v1或v2。 - 校验自定义 state schema 至少有
messages、remaining_steps;有 structured response 时还要有structured_response。 - 没有 schema 时使用
AgentState或AgentStateWithStructuredResponse。 - 将传入 tools 分成普通工具和 provider builtin tool schema,构造或复用
ToolNode。 - 判断 model 是静态模型、Runnable,还是运行时动态选择模型。
- 静态模型若是字符串,调用
langchain.chat_models.init_chat_model()初始化 provider。 _should_bind_tools()检查模型是否已有 tools binding;需要时执行model.bind_tools(tool_classes + llm_builtin_tools)。- 构造
static_model = _get_prompt_runnable(prompt) | model。
这里的 | 不是普通 Python 管道,而是构造一个 Runnable sequence:第一步把 graph state 转成 LLM input,第二步把 input 交给 ChatModel。
6.2 _get_prompt_runnable():四种 prompt 形态
源码分支对应的实际输入输出:
| prompt 类型 | 构造的 Runnable | LLM 最终收到的内容 |
|---|---|---|
None |
lambda state: state["messages"] |
原始消息列表 |
str |
创建 SystemMessage(content=prompt),再拼接 |
[system_message] + state["messages"] |
SystemMessage |
把该消息放在列表首位 | [prompt] + state["messages"] |
| callable | RunnableCallable(prompt) |
callable(state) 的结果 |
| coroutine callable | RunnableCallable(None, prompt) |
await prompt(state) 的结果 |
| Runnable | 原样使用 | Runnable 的输出 |
注意两点:
- 字符串 prompt 不是拼进用户消息的字符串,而是转换成单独的
SystemMessage。 - prompt callable 拿到的是完整 graph state,不只是 messages,因此可以根据用户、租户、检索结果或运行上下文动态生成 LLM input。
6.3 RunnableSeq.invoke():prompt 和 model 的真实调用顺序
内部实现:libs/langgraph/langgraph/_internal/_runnable.py。
RunnableSeq.invoke(input, config):
- 创建 callback manager 和根 run。
- 对
steps逐个循环。 - 第一个 step 是 Prompt Runnable,在 config context 中执行。
- 第一个 step 的返回值成为第二个 step 的 input。
- 第二个 step 是 ChatModel 的
invoke(model_input, config)。 - 正常结束调用
on_chain_end;异常调用on_chain_error。
于是静态 LLM 的确切调用链是:
1 | |
LangGraph 不负责把 token 生成出来;模型 provider 负责生成。LangGraph 通过 callback 和 StreamMessagesHandler 观察模型事件,再把 token/metadata 放入 graph stream。
6.4 call_model():每轮 agent 节点做什么
同步路径源码顺序:
- 如果传入 async dynamic model 却调用同步 agent,立即报错。
_get_model_input_state(state)取得本轮 LLM 输入消息。- 有
pre_model_hook时优先取llm_input_messages,否则取 state 的messages。 - 没有消息则抛出明确的
ValueError;有消息则_validate_chat_history()校验消息序列。 - 把选出的 messages 临时放回
state["messages"],因为 Prompt Runnable 约定从这个 key 读取。 - 静态模型直接
static_model.invoke(model_input, config);动态模型先_resolve_model(state, runtime),再 invoke。 - 将返回的
AIMessage.name设为 agent name。 _are_more_steps_needed()检查 remaining_steps 和 tool calls。- 如果步骤不足,返回一条内容为
Sorry, need more steps to process this request.的 AIMessage。 - 否则返回
{"messages": [response]}。返回 list 是为了让 messages channel 的 reducer 追加消息,而不是替换历史。
异步路径完全对应,只将模型调用替换为 await dynamic_model.ainvoke(...) 或 await static_model.ainvoke(...)。
6.5 pre-model hook 如何改变 prompt
hook 的输出可以包含:
1 | |
含义不同:
messages会更新持久状态,适合摘要、裁剪和重写历史。llm_input_messages只作为本次模型输入,不写回 messages。- 其他 key 继续进入 graph state。
因此”prompt 拼接”发生在两个可能位置:先是 pre-model hook 产生模型输入,再是 _get_prompt_runnable() 在输入前插入 system message 或执行自定义 prompt Runnable。
7. ReAct Agent 的完整 loop:LLM、工具和回边
7.1 图结构
经典逻辑可表示为:
flowchart TD
START --> PRE[pre_model_hook 可选]
PRE --> AGENT[agent / call_model]
START --> AGENT
AGENT --> ROUTE{AIMessage.tool_calls?}
ROUTE -- 无 --> STRUCT[structured response 可选]
ROUTE -- 有 --> TOOLS[ToolNode]
TOOLS --> AGENT
STRUCT --> END
真正的循环不是 call_model() 内部的 while:
1 | |
每次回到 agent 都会重新执行 prompt Runnable 和 ChatModel,因此模型每轮看到的是完整的消息历史(或 hook 指定的裁剪历史)。
7.2 条件路由如何决定停止或回 tools
tools_condition 检查最后一条消息是否是带 tool_calls 的 AIMessage:
- 有 tool calls:目标为
tools。 - 没有 tool calls:目标为
END或 structured response 节点。
return_direct 工具和 remaining_steps 会额外影响是否允许继续。remaining steps 小于 2 且仍有 tool calls 时,不再调工具,而是由 call_model() 返回”步骤不足”的最终 AIMessage,避免直接触发 GraphRecursionError。
7.3 ToolNode 如何执行一个 tool call
文件:libs/prebuilt/langgraph/prebuilt/tool_node.py。
ToolNode._func() 先解析输入形态:可以是 message list、带 messages 的 state dict,或直接 tool calls。然后对 tool calls 使用 executor 并行执行 _run_one()。
单个调用的细流程:
- 从 call 取
name、args、id,在tools_by_name中查找工具。 - 找不到工具时
_validate_tool_call()返回错误ToolMessage,内容列出可用工具名。 _inject_tool_args()在真正调用前注入 state、store、runtime 等特殊参数;这些参数不会来自 LLM 的 args。- 调用
tool.invoke(call_args, config)。 ValidationError被包装为ToolInvocationError,错误报告使用原始 args,避免把注入参数暴露给模型。_normalize_tool_response()把字符串、消息、Command 或列表统一转换为ToolMessage/Command writes。GraphBubbleUp(包括 interrupt)直接向上抛出,不能被普通 tool error handler 吞掉。- 其他异常按
handle_tool_errors策略处理;默认转换为 status=error的ToolMessage。 - 多个 tool call 的结果合并成
{"messages": [ToolMessage, ...]},由 messages channel 追加到历史。
7.4 v1 与 v2 的差异
| 版本 | 一个 ToolNode task 处理什么 | 并发位置 |
|---|---|---|
| v1 | 最后一条 AIMessage 中的全部 tool calls | ToolNode 内部 executor 并行 |
| v2 | 一个 tool call | 图通过 Send 为每个 call 创建独立 ToolNode task |
两种版本最终都走同一个 Pregel loop、checkpoint、retry、interrupt 和 stream 机制。v2 的优势是每个工具调用有独立 task identity、路径和流事件,更适合细粒度分发和嵌套图。
7.5 structured response 是额外 LLM 调用
当提供 response_format 时,Agent loop 停止后才执行 generate_structured_response:
- 解析 schema 或
(prompt, schema)二元组。 - 对当前模型调用
.with_structured_output(schema)。 - 用最终消息历史(以及可选的 structured prompt)再次调用 LLM。
- 把结果写入
structured_response。
它不是同一次 tool-calling response 的自动字段,也不是免费解析;源码文档明确说明这是 Agent loop 结束后的单独模型调用。
8. 状态、消息和 reducer 的真实语义
8.1 messages 为什么能累加
Agent state 的 messages 通常不是普通 LastValue,而是带消息 reducer 的 channel。节点返回:
1 | |
不会把整个历史替换成一条消息,而是由 channel reducer 追加/合并;ToolMessage、AIMessage、RemoveMessage 也由该 reducer 处理。若 pre-model hook 要重写历史,必须显式发送 RemoveMessage(id=REMOVE_ALL_MESSAGES),否则只是追加而不是覆盖。
8.2 state、context、store、checkpoint 的边界
| 对象 | 生命周期 | 用途 |
|---|---|---|
| state/channels | 图执行期间,按 step 演进 | 当前会话和节点间数据 |
| context | 一次 run | 模型选择、租户、用户配置等只读运行上下文 |
| checkpoint | thread + namespace 的时间线 | 恢复、中断、回放、时间旅行 |
| store | 跨 thread | 用户记忆、共享知识、应用数据 |
不要用 store 代替 checkpoint:store 保存业务数据,不保存可恢复的 Pregel task/channel 时间线。
9. Checkpoint 与恢复
基础契约:libs/checkpoint/langgraph/checkpoint/base/__init__.py。
关键数据:
Checkpoint:channel values、channel versions、versions_seen、id、timestamp 等。CheckpointTuple:checkpoint、metadata、父 config 和 pending writes。BaseCheckpointSaver:get_tuple/list/put/put_writes/delete_thread以及异步对应方法。
运行身份主要来自:
1 | |
一次可恢复运行包含两类写入:
put_writes:节点完成后立即保存 task 级 pending writes,支持节点执行到一半发生中断/失败后恢复。put:after_tick()在超步边界保存聚合后的 checkpoint。
恢复时 tick() 会把已经成功保存的普通 writes 重新挂回 task,并跳过已成功节点;控制类 writes(ERROR、INTERRUPT、RESUME)不能按普通状态写入重放。这样恢复不会简单地从头重跑所有节点,也不会重复追加已经持久化的消息。
SQLite 的表结构可以作为最直观的后端参照:checkpoints 保存快照和 metadata,writes 保存 task/channel/value。Postgres 实现遵守同一契约,改变的是存储并发和事务方式,不改变图执行语义。
10. Stream、LLM token 和事件从哪里来
Pregel.stream() 自身不是简单 yield final_state。事件来源有三条:
- loop 在
after_tick()等位置发出的 values/updates/checkpoint/debug。 - LangChain callback handler 观察 ChatModel 的 token、tool call 和 metadata,写入
SyncQueue。 - 节点通过 Runtime 的
stream_writer写入 custom 数据。
当 stream_mode="messages" 时,LLM 调用链中的 callback 事件会被 StreamMessagesHandler 捕获;因此模型 token 可以在一个 Pregel task 尚未结束时被输出。updates 则更接近节点提交结果,通常要等 task 完成后才能看到。
嵌套图通过 checkpoint_ns 和 namespace tuple 区分事件;subgraphs=True 时事件会带上父节点和子节点的路径,避免不同子图的同名 agent/tool 事件混淆。
11. 同步、异步、重试和并发的边界
stream/invoke走SyncPregelLoop、同步 runner 和Runnable.invoke。astream/ainvoke走异步 loop、PregelRunner.atick和Runnable.ainvoke。- 一个超步内无依赖的 tasks 可以并发执行;下一个超步必须等待当前超步的 writes 应用。
- retry policy 包裹的是 task attempt,不是整个 graph;失败重试时不会把成功的其他 task 重新执行。
- 节点 error handler 是 graph task,普通错误被路由到 handler;handler 自身失败不能再次捕获自己。
GraphInterrupt/GraphBubbleUp是控制流异常,不应被普通工具错误处理吞掉。
12. 按问题定位的源码路径
| 问题 | 第一跳 | 第二跳 | 重点看什么 |
|---|---|---|---|
| 节点没有执行 | StateGraph.compile |
_algo.prepare_next_tasks |
channel trigger、START、条件边 |
| 节点执行顺序不对 | Pregel.stream |
SyncPregelLoop.tick/after_tick |
BSP 边界、updated channels |
| state 被覆盖 | _add_schema |
对应 channel 的 update |
LastValue、reducer、messages channel |
| LLM 没看到 system prompt | _get_prompt_runnable |
RunnableSeq.invoke |
prompt 类型、model_input |
| LLM 被调用多次 | create_react_agent.call_model |
tools_condition |
是否有 tool_calls、是否回 agent |
| 工具没执行 | ToolNode._func |
_validate_tool_call、_run_one |
工具名、注入参数、工具错误策略 |
| Agent 无限循环 | call_model |
_are_more_steps_needed、remaining_steps |
模型持续返回 tool_calls、路由是否到 END |
| 中断无法恢复 | SyncPregelLoop |
put_writes、checkpoint saver |
thread_id、checkpoint_ns、pending writes |
| stream 没 token | Pregel.stream |
StreamMessagesHandler |
stream mode、callback 是否传递 |
| 并行结果不稳定 | _runner.PregelRunner.tick |
_algo.apply_writes |
reducer 是否可交换/可结合、超步边界 |
13. 推荐的阅读顺序
libs/cli/uv-examples/simple/src/agent/graph.py:看最小业务图如何声明。libs/langgraph/langgraph/graph/state.py:读StateGraph.__init__、_add_schema、具体add_node、add_edge、compile、CompiledStateGraph.attach_*。libs/langgraph/langgraph/pregel/main.py:读invoke、stream的配置、handler、loop、runner 和异常收尾。libs/langgraph/langgraph/pregel/_loop.py:逐行读__init__、tick、after_tick、put_writes。libs/langgraph/langgraph/pregel/_algo.py:读prepare_next_tasks、apply_writes,理解 channel version 和 trigger。libs/langgraph/langgraph/pregel/_runner.py:读PregelRunner.tick、commit、error handler 和 retry 交界。libs/langgraph/langgraph/_internal/_runnable.py:读RunnableCallable、RunnableSeq.invoke/ainvoke,确认 prompt/model 的实际调用顺序。libs/prebuilt/langgraph/prebuilt/chat_agent_executor.py:读_get_prompt_runnable、create_react_agent、call_model、generate_structured_response。libs/prebuilt/langgraph/prebuilt/tool_node.py:读_func、_run_one、_execute_tool_sync、_inject_tool_args、_normalize_tool_response。- 最后读 checkpoint base 和一个具体后端,理解 pending writes 如何持久化和恢复。
14. 关键符号索引
| 符号 | 文件 | 作用 |
|---|---|---|
StateGraph |
libs/langgraph/langgraph/graph/state.py |
声明层 builder |
StateGraph._add_schema |
同上 | schema -> channel/managed |
StateGraph.compile |
同上 | builder -> Pregel 图 |
CompiledStateGraph.attach_node |
同上 | 用户节点 -> PregelNode |
Pregel.invoke |
libs/langgraph/langgraph/pregel/main.py |
收集运行结果 |
Pregel.stream |
同上 | 配置、调度、事件输出主入口 |
SyncPregelLoop.tick |
libs/langgraph/langgraph/pregel/_loop.py |
准备一个超步 |
SyncPregelLoop.after_tick |
同上 | 应用 writes、checkpoint、interrupt |
prepare_next_tasks |
libs/langgraph/langgraph/pregel/_algo.py |
channel trigger -> tasks |
apply_writes |
同上 | task writes -> channel values |
PregelRunner.tick |
libs/langgraph/langgraph/pregel/_runner.py |
并发执行任务 |
RunnableSeq.invoke |
libs/langgraph/langgraph/_internal/_runnable.py |
prompt -> model 串联 |
_get_prompt_runnable |
libs/prebuilt/langgraph/prebuilt/chat_agent_executor.py |
prompt 规范化和消息拼接 |
call_model |
同上 | 取 state、调用 LLM、写回 AIMessage |
ToolNode._func |
libs/prebuilt/langgraph/prebuilt/tool_node.py |
解析并执行 tool calls |
ToolNode._execute_tool_sync |
同上 | 注入参数、调用工具、转 ToolMessage |
BaseCheckpointSaver |
libs/checkpoint/langgraph/checkpoint/base/__init__.py |
checkpoint 后端契约 |
15. 结语
把全文压缩成一句话:LangGraph 把业务函数、LLM 和工具都统一成可读写 channel 的 Runnable task;Pregel loop 以超步为边界执行、合并、持久化和路由,Agent loop 只是由 AIMessage 的 tool call 触发的图回边。
理解了这一点,再回头看 create_react_agent、ToolNode、stream_mode="messages" 这些 API,它们就不再是黑盒,而是同一套运行时机制在不同层面的投影。