LangGraph
前面我们介绍了第一个 Agent Framework - Pocket Flow。并且说道,大多数 Agent Framework 的核心是提供了对 Workflow 的抽象-Graph。今天我们来介绍另一个 Agent Framework - LangGraph。并依旧重点关注 Graph 实现的三个问题:
- 如何表示 Graph 中的节点以及节点的触发关系
- 如何在节点之间传递共享数据
- Graph 如何被驱动执行
Langgraph 比 Pocket Flow 代码复杂的多,主要有如下几个原因:
- Langgraph 定义的 Graph 所能表达的语义更加丰富。在 Pocket Flow 中,一个 action 只能触发一个节点,Langgraph 中一次可以触发多个节点,一节点也可以定义对多个节点的依赖。
- Langgraph 支持节点的并发执行,功能也更加完善
Langgraph Graph 分为定义和运行时两种表示。我们将从一个示例开发,直接看 Graph 的运行时表示,这样可以更清晰的解释我们所关心的三个问题。
1. 从一个完整的 Graph 开始
我们先定义一个文章生成流程:先生成提纲,再由两个节点并行写出候选稿;两个节点都完成后进入审核;审核通过则结束,否则修改后再次审核。
|
|
这个例子包含了顺序执行、并行、汇合、条件分支和循环。
|
|
从使用者的角度看,Node 是计算逻辑,Edge 是节点之间的触发关系,State 是节点之间的共享数据。但 StateGraph 只是 Builder,调用 compile() 后,它会被转换成另一组运行时结构,并由 Pregel 执行。
2. Graph 运行时包含哪些数据结构
builder 是 StateGraph,保存用户定义的 Graph:
|
|
graph 是 CompiledStateGraph。它继承自 Pregel,保存真正执行 Graph 所需的结构:
|
|
运行时的核心是:PregelNode 通过 channels 读取数据,通过 writers 写数据和控制信号;Pregel 根据控制 Channel 的更新匹配节点的 triggers。
在分别回答三个问题之前,需要先理解这个匹配发生在哪个时间点。
3. Pregel 的三阶段执行
Pregel 把 Graph 拆成一轮一轮的 Superstep。每轮固定分成三个阶段:
|
|
关键点是:本轮 writer 写出的控制 Channel,要到 Update 阶段才真正更新;订阅它的 trigger,要到下一轮 Plan 才会被匹配。
以 plan -> writer_a 为例。编译后有下面两个接口:
|
|
它们在相邻两轮之间这样连接:
|
|
3.1 先区分“值的可见性”和“更新是否已处理”
这里其实有两个不同的问题:
- 本轮节点写出的值,什么时候能被其他节点读取?
- 一个 Channel 中存在值时,如何判断它是不是节点尚未处理的新更新?
第一个问题由 Pregel 的阶段边界解决,第二个问题才由 Channel Version 解决。
值的可见性:在 Update 阶段统一提交
假设 writer_a 和 writer_b 在同一个 Superstep 中并行执行。两者开始执行时,读取的是同一份 State:
|
|
执行过程中,它们分别产生 Write:
|
|
这些 Write 先保存在各自的 Task 中,不会立即修改 drafts Channel。因此在本轮 Execution 阶段,两个节点都看不到对方的结果。
等两个节点全部完成,Pregel 才在 Update 阶段统一提交:
|
|
下一轮执行的 review 才能读到这个新值。
所以:
“本轮写、下轮见”是由 Pending Writes 延迟到 Update 阶段统一提交实现的,不是 Version 实现的。
Channel Version:判断更新是不是新的
Version 是怎么生成的
版本保存在:
|
|
每次进入 Update 阶段,apply_writes() 先找出所有 Channel 中的最大版本,再调用一次版本生成函数:
|
|
如果没有配置 Checkpointer,get_next_version 就是一个简单的自增函数:
|
|
如果配置了 Checkpointer,则使用 checkpointer.get_next_version()。默认实现也是加一;具体实现可以使用 int、float 或 str,只要生成的版本单调递增、能够比较大小。比如 InMemorySaver 使用以递增整数开头的字符串版本。
同一个 Superstep 只生成一次 next_version。本轮所有真正发生变化的 Channel 都被赋予这个版本;没有变化的 Channel 保留原版本。因此它更像一个 Superstep 的逻辑时钟,而不是每个 Channel 各自独立的计数器。
例如 Update 前:
|
|
当前最大版本是 2,所以本轮生成 next_version = 3。假设本轮更新了 drafts 和 Barrier Channel,Update 后就是:
|
|
版本号本身没有业务含义。Pregel 只关心大小关系:当前版本是否大于节点记录在 versions_seen 中的版本。
接着只看 plan -> writer_a 这条边。编译后,两端通过同一个控制 Channel 连接:
|
|
Pregel 为每个 Channel 记录当前版本,同时记录每个节点上次处理到的版本:
|
|
假设一开始:
|
|
plan.writers 在本轮写入这个控制 Channel。Update 阶段提交后,Channel 获得一个新版本:
|
|
下一轮 Plan 检查 writer_a.triggers:
|
|
writer_a 处理完这次触发后,Pregel 记录:
|
|
现在两个版本相等,这次更新已经被处理,不会再次触发:
|
|
如果以后 Graph 再次执行到 plan,它又向同一个控制 Channel 写入信号,Channel 会产生一个更新的版本,例如 v2:
|
|
实际判断还要求 Channel 当前可用。对于 branch:to:* 使用的 EphemeralValue,信号消费后还会被清空;版本比较提供了一套适用于所有 Channel 的统一“新旧判断”。
因此 Version 解决的是:
不要问 Channel 里有没有值,而要问这个节点有没有处理过 Channel 的当前版本。
updated_channels 是什么
Update 阶段结束后,apply_writes() 返回一个 set[str],记录本轮成功更新且当前可用的 Channel 名称。例如 plan 执行后:
|
|
它只包含 Channel 名称,不包含值,也不直接包含节点。下一轮 Plan 用它查询反向索引:
|
|
匹配结果是:
|
|
这里有一个容易误解的地方:既然 Channel 刚刚出现在 updated_channels 中,它的版本必然刚刚增大,为什么还要和 versions_seen 比较?
答案是:在一次普通、连续执行的快速路径中,这个版本检查通常确实必然通过。
例如:
|
|
在这个限定场景里,找到 writer_a 后直接创建 Task,结果也是一样的。源码仍然比较版本,是因为两者承担的职责不同:
updated_channels是全局的变化集合,用来快速缩小候选节点范围versions_seen[node][channel]是每个节点自己的消费游标,用来判断该节点是否处理过当前版本
prepare_next_tasks() 不只服务于连续执行的快速路径,还要处理 checkpoint 恢复、重放以及没有 updated_channels 的情况。例如恢复时可能只有:
|
|
虽然没有 updated_channels 可以提供候选集合,Pregel 扫描节点后仍能通过 5 > 4 判断 writer_a 尚未处理这次更新。如果两边都是 5,则说明已经处理过,恢复时不能重复执行。
因此,updated_channels 不是调度的权威状态,而是一个可选的索引优化;Channel Version 和 versions_seen 才是能够跨 checkpoint 恢复的最终判断依据。
整个过程可以浓缩成:
|
|
下面再分别回答开头的三个问题。
4. 如何表示节点以及节点的触发关系
4.1 PregelNode 的四个核心接口
add_node() 先把函数保存成 StateNodeSpec。编译后,它会变成 PregelNode:
|
|
四个字段组成一次完整的节点执行:
|
|
triggers 和 writers 都保存 Channel 层面的接口:
triggers是 Channel 名称列表。任意一个 Channel 出现节点尚未消费的新版本,节点就可以被触发writers是节点函数执行后的 Runnable 列表。它们把返回值写入数据 Channel,也把出边转换成对控制 Channel 的写入
所以可以把 LangGraph 的触发关系理解成:源节点的 writers 和目标节点的 triggers 通过同名控制 Channel 连接。
4.2 普通边
|
|
attach_edge() 为这条边向 plan.writers 追加一个 ChannelWrite:
|
|
而 writer_a 的 PregelNode 订阅同名 Channel:
|
|
writer 只产生 Write,不直接运行目标节点。Write 在本轮 Update 阶段更新 Channel,下一轮 Plan 才通过 trigger_to_nodes 和 Channel 版本找到 writer_a。
4.3 并行、汇合与条件边
plan 有两条出边,因此它的 writers 会同时写入两个控制 Channel,两个 writer 在下一轮并行执行。
|
|
多起点边会转换成 NamedBarrierValue Channel:
|
|
两个 writer 分别向 Barrier 写入自己的名字。只有 Barrier 收齐两个值并变为可用,它的版本更新才会在下一轮触发 review。
条件边则由路由函数决定 writer 最终写入哪个控制 Channel:
|
|
route 返回 revise 时写入 branch:to:revise;返回 END 时结束。无论普通边、汇合边还是条件边,最后都统一成 writer、控制 Channel 和 trigger 的关系。
5. 如何在节点之间传递共享数据
API 层的节点共享一个 State:
|
|
运行时并不存在一个由所有节点直接修改的全局字典。LangGraph 会把 State 的每个字段转换成一个 Channel:
|
|
PregelNode.channels 指定节点读取哪些数据 Channel。节点执行前,LangGraph 读取这些 Channel 的当前值,组合成 State;节点返回的 Partial State 再由 writers 拆成数据 Channel Writes。
例如 plan 返回:
|
|
表示向 outline Channel 提交一次更新,而不是直接修改全局字典。
普通字段默认使用 LastValue,一个 Superstep 中最多接收一个更新。writer_a 和 writer_b 会在同一轮写入 drafts,所以它声明了 reducer:
|
|
它会转换成 BinaryOperatorAggregate:
|
|
因此,State 是节点看到的共享数据视图,Channel 才是共享数据实际的保存和更新机制。
6. Graph 被驱动执行的过程
现在可以把执行过程完整串起来:
|
|
源码的主循环可以简化为:
|
|
当 Update 阶段产生的 Channel 不再匹配任何节点的 triggers 时,下一轮 Plan 得不到 Task,Graph 执行结束。
7. 用户代码如何转换成运行时结构
compile() 把用户定义的 Workflow 转换成 Channel 网络:
|
|
具体的转换关系如下:
| 用户代码 | StateGraph 中的结构 |
编译后的运行时结构 |
|---|---|---|
State 字段 |
channels |
LastValue 或带 reducer 的数据 Channel |
add_node() |
StateNodeSpec |
包含 channels/triggers/bound/writers 的 PregelNode |
add_edge(A, B) |
edges |
A 的 writer 与 B 的 trigger 连接同一个控制 Channel |
add_edge([A, B], C) |
waiting_edges |
writers 写入、C 订阅的 Barrier Channel |
add_conditional_edges() |
BranchSpec |
根据路由结果写入目标控制 Channel |
到这里,开头的三个问题可以统一到同一个模型中:
- 节点由
PregelNode表示,触发关系由 writers、控制 Channel 和 triggers 表示 - State 的每个字段对应一个数据 Channel,节点通过读取 Channel 和提交 Channel Write 共享数据
- Pregel 按照
Plan -> Execution -> Update推进 Graph,并在相邻轮次之间把 writer 的输出匹配到 trigger
Node 负责计算,writer 负责写 Channel,trigger 负责订阅 Channel,Pregel 负责在相邻轮次之间完成匹配和调度。
8. 源码阅读
本文基于 LangGraph 1.2.10、commit 658541c4960f329864a2523fc7d52427e8190bed: