Pregel 实现了 LangGraph 的运行时,负责管理 LangGraph 应用程序的执行。
编译 StateGraph 或创建 @entrypoint 会生成一个 Pregel 实例,该实例可以通过输入进行调用。
本指南从较高层面解释运行时,并提供使用 Pregel 直接实现应用程序的说明。
注意: Pregel 运行时得名于 Google 的 Pregel 算法,该算法描述了一种使用图进行大规模并行计算的高效方法。
在 LangGraph 中,Pregel 将 actor 和 channel 组合成一个应用程序。Actor 从 channel 中读取数据,并向 channel 写入数据。Pregel 按照 Pregel Algorithm/Bulk Synchronous Parallel 模型,将应用程序的执行组织为多个步骤。
每个步骤包含三个阶段:
- 计划:确定此步骤中要执行哪些 actor。例如,在第一个步骤中,选择订阅特殊 input channel 的 actor;在后续步骤中,选择订阅上一步骤中已更新 channel 的 actor。
- 执行:并行执行所有被选中的 actor,直到全部完成、某个 actor 失败,或达到超时时间。在此阶段,channel 的更新对 actor 不可见,直到下一个步骤才可见。
- 更新:使用此步骤中 actor 写入的值更新 channel。
重复执行,直到没有 actor 被选中执行,或达到最大步骤数。
Actor
actor 是一个 PregelNode。它订阅 channel,从中读取数据,并向其中写入数据。可以将它理解为 Pregel 算法中的一个 actor。PregelNodes 实现了 LangChain 的 Runnable 接口。
Channel
Channel 用于 actor(PregelNodes)之间的通信。每个 channel 都有一个值类型、一个更新类型,以及一个更新函数——该函数接收一系列更新并修改已存储的值。Channel 可用于将数据从一条 chain 发送到另一条 chain,或在未来某个步骤中将数据从一条 chain 发送回自身。
LastValue
LastValue 是默认的 channel 类型。它会存储最后一次写入的值,并覆盖任何先前的值。可将它用于输入和输出值,或用于将数据从一个步骤传递到下一个步骤。
from langgraph.channels import LastValue
channel: LastValue[int] = LastValue(int)
Topic
Topic 是一个可配置的 PubSub channel,适合在 actor 之间发送多个值,或跨步骤累积输出。它可以配置为对值去重,或累积一次运行期间写入的所有值。
from langgraph.channels import Topic
# Accumulate all values written across steps
channel: Topic[str] = Topic(str, accumulate=True)
BinaryOperatorAggregate
BinaryOperatorAggregate 存储一个持久值,并通过将二元运算符应用于当前值和每个新更新来更新该值。可使用它跨步骤计算运行中的聚合结果。
import operator
from langgraph.channels import BinaryOperatorAggregate
# Running total: each write adds to the current value
total = BinaryOperatorAggregate(int, operator.add)
DeltaChannel(beta)
DeltaChannel 需要 langgraph>=1.2,目前处于 beta 阶段。API 可能会在未来版本中发生变化。
DeltaChannel 在每个步骤中只存储增量 delta,而不是完整的累积值。这对于写入频繁、并且随着时间不断累积大量值的 channel 最有用——例如,长时间运行线程中的对话消息列表。如果不使用 delta 存储,完整列表会被重新序列化到每个 checkpoint 中;使用 DeltaChannel 后,只会存储每个步骤中新写入的消息。
当某个 channel 写入频繁,并且会随着时间增长得很大时,可以考虑使用 DeltaChannel。一个很好的判断信号是:如果你发现某个特定 channel 的 checkpoint 大小会随着 thread 长度线性增长,那么 DeltaChannel 很可能很适合。
像使用普通 reducer 一样,在 Annotated 类型注解中使用 DeltaChannel:
from typing import Annotated, Sequence
from typing_extensions import TypedDict
from langgraph.channels import DeltaChannel
def my_reducer(state: list[str], writes: Sequence[list[str]]) -> list[str]:
result = list(state)
for write in writes:
result.extend(write)
return result
class State(TypedDict):
messages: Annotated[list[str], DeltaChannel(my_reducer)]
批量 reducer 要求
传给 DeltaChannel 的 reducer 是一个 批量 reducer:它会在一次调用中接收当前状态和当前步骤中所有写入组成的 sequence,而不是像标准 reducer 那样成对处理。这与 StateGraph 中配合 Annotated 使用的按键 reducer 不同,后者会针对每次更新调用一次 reducer。
批量 reducer 必须是结合的(与批处理方式无关):reducer(reducer(state, [xs]), [ys]) == reducer(state, [xs, ys])
如果你的 reducer 不满足结合性,重建出的状态可能会因 LangGraph 跨步骤对写入进行批处理的方式不同而不同,从而产生不一致的行为。
reducer 在重建时运行,而不是在写入时运行。 与 BinaryOperatorAggregate 不同,后者的 reducer 会在写入时调用,因此被序列化到 checkpoint 中的是合并后的值;而 DeltaChannel 的 reducer 会在 channel 值从持久化写入中被重建时调用。被序列化的是每个步骤的原始写入;reducer 只会在值被物化时调用——例如下一次读取时、下一步骤的 actor 运行时,或重放历史时。设计 reducer 时的实际影响:
- 将它设计为
(state, writes) 的纯函数。 任何副作用、随机性或读取当前时间的行为(例如 uuid.uuid4()、datetime.now())都会在每次值被重建时执行,并在每次重放时产生不同结果。它们不会被固化到持久化写入中。
- 不要依赖对传入写入的修改会被持久化。 如果你的 reducer 修改了某个写入对象(例如,为一个没有 ID 的条目分配稳定 ID),这个修改只存在于重建后的值中。存储的写入仍然保持原始形态,因此下一次重建时仍会看到未修改的输入。
- 在上游附加标识和其他稳定元数据。 如果下游代码需要跨轮次通过 ID 引用某个条目(例如之后更新或删除它),请在将值写入 channel 之前分配该 ID,而不是在 reducer 内部分配。
以下是两种最常见场景的批量 reducer:
from typing import Any, Sequence
# List: append all writes in order
def list_reducer(state: list[Any], writes: Sequence[list[Any]]) -> list[Any]:
result = list(state)
for write in writes:
result.extend(write)
return result
# Dict: merge all writes, last write wins on key conflicts
def dict_reducer(
state: dict[str, Any], writes: Sequence[dict[str, Any]]
) -> dict[str, Any]:
result = dict(state)
for write in writes:
result.update(write)
return result
二者都满足结合性:逐个应用批次与一次性一起应用会得到相同结果。
使用 snapshot_frequency 限制读取延迟
如果没有 snapshot,读取 DeltaChannel 值需要重放完整的写入历史——对于包含 N 个步骤的 thread 来说是 O(N)。设置 snapshot_frequency=K 会每 K 个 pregel 步骤写入一个完整 snapshot,将读取深度限制在最多 K 个步骤:
class State(TypedDict):
messages: Annotated[
list[str],
DeltaChannel(my_reducer, snapshot_frequency=5),
]
较高的 snapshot_frequency 值会降低存储开销,但会增加读取延迟。较低的值能更严格地限制延迟,但代价是 checkpoint 更大。None(默认值)会完全跳过 snapshot——适用于读取较少或 thread 较短的情况。
虽然大多数用户会通过 StateGraph API 或 @entrypoint 装饰器与 Pregel 交互,但也可以直接与 Pregel 交互。
下面给出几个不同的示例,帮助你了解 Pregel API。
Single node
Multiple nodes
Topic
BinaryOperatorAggregate
Cycle
from langgraph.channels import EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b")
)
app = Pregel(
nodes={"node1": node1},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
},
input_channels=["a"],
output_channels=["b"],
)
app.invoke({"a": "foo"})
from langgraph.channels import LastValue, EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b")
)
node2 = (
NodeBuilder().subscribe_only("b")
.do(lambda x: x + x)
.write_to("c")
)
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": LastValue(str),
"c": EphemeralValue(str),
},
input_channels=["a"],
output_channels=["b", "c"],
)
app.invoke({"a": "foo"})
{'b': 'foofoo', 'c': 'foofoofoofoo'}
from langgraph.channels import EphemeralValue, Topic
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b", "c")
)
node2 = (
NodeBuilder().subscribe_to("b")
.do(lambda x: x["b"] + x["b"])
.write_to("c")
)
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
"c": Topic(str, accumulate=True),
},
input_channels=["a"],
output_channels=["c"],
)
app.invoke({"a": "foo"})
{'c': ['foofoo', 'foofoofoofoo']}
此示例演示如何使用 BinaryOperatorAggregate channel 来实现 reducer。from langgraph.channels import EphemeralValue, BinaryOperatorAggregate
from langgraph.pregel import Pregel, NodeBuilder
node1 = (
NodeBuilder().subscribe_only("a")
.do(lambda x: x + x)
.write_to("b", "c")
)
node2 = (
NodeBuilder().subscribe_only("b")
.do(lambda x: x + x)
.write_to("c")
)
def reducer(current, update):
if current:
return current + " | " + update
else:
return update
app = Pregel(
nodes={"node1": node1, "node2": node2},
channels={
"a": EphemeralValue(str),
"b": EphemeralValue(str),
"c": BinaryOperatorAggregate(str, operator=reducer),
},
input_channels=["a"],
output_channels=["c"],
)
app.invoke({"a": "foo"})
{ 'c': 'foofoo | foofoofoofoo' }
此示例演示如何在图中引入一个循环,方式是让
一条 chain 向它所订阅的 channel 写入数据。执行会持续进行,
直到向该 channel 写入一个 None 值。from langgraph.channels import EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder, ChannelWriteEntry
example_node = (
NodeBuilder().subscribe_only("value")
.do(lambda x: x + x if len(x) < 10 else None)
.write_to(ChannelWriteEntry("value", skip_none=True))
)
app = Pregel(
nodes={"example_node": example_node},
channels={
"value": EphemeralValue(str),
},
input_channels=["value"],
output_channels=["value"],
)
app.invoke({"value": "a"})
{'value': 'aaaaaaaaaaaaaaaa'}
高级 API
LangGraph 提供了两个用于创建 Pregel 应用程序的高级 API:StateGraph (Graph API) 和 Functional API。
StateGraph (Graph API)
Functional API
StateGraph (Graph API) 是一种更高层的抽象,可简化 Pregel 应用程序的创建。它允许你定义由节点和边组成的图。编译该图时,StateGraph API 会自动为你创建 Pregel 应用程序。from typing import TypedDict
from langgraph.constants import START
from langgraph.graph import StateGraph
class Essay(TypedDict):
topic: str
content: str | None
score: float | None
def write_essay(essay: Essay):
return {
"content": f"Essay about {essay['topic']}",
}
def score_essay(essay: Essay):
return {
"score": 10
}
builder = StateGraph(Essay)
builder.add_node(write_essay)
builder.add_node(score_essay)
builder.add_edge(START, "write_essay")
builder.add_edge("write_essay", "score_essay")
# Compile the graph.
# This will return a Pregel instance.
graph = builder.compile()
编译后的 Pregel 实例会关联一组节点和 channel。你可以通过打印它们来查看这些节点和 channel。你会看到类似如下内容:{'__start__': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1810>,
'write_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba14d0>,
'score_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1710>}
你应该会看到类似如下内容{'topic': <langgraph.channels.last_value.LastValue at 0x7d05e3294d80>,
'content': <langgraph.channels.last_value.LastValue at 0x7d05e3295040>,
'score': <langgraph.channels.last_value.LastValue at 0x7d05e3295980>,
'__start__': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e3297e00>,
'write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e32960c0>,
'score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8ab80>,
'branch:__start__:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e32941c0>,
'branch:__start__:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d88800>,
'branch:write_essay:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e3295ec0>,
'branch:write_essay:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8ac00>,
'branch:score_essay:__self__:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d89700>,
'branch:score_essay:__self__:score_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8b400>,
'start:write_essay': <langgraph.channels.ephemeral_value.EphemeralValue at 0x7d05e2d8b280>}
在 Functional API 中,你可以使用 @entrypoint 创建 Pregel 应用程序。entrypoint 装饰器允许你定义一个接收输入并返回输出的函数。from typing import TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.func import entrypoint
class Essay(TypedDict):
topic: str
content: str | None
score: float | None
checkpointer = InMemorySaver()
@entrypoint(checkpointer=checkpointer)
def write_essay(essay: Essay):
return {
"content": f"Essay about {essay['topic']}",
}
print("Nodes: ")
print(write_essay.nodes)
print("Channels: ")
print(write_essay.channels)
Nodes:
{'write_essay': <langgraph.pregel.read.PregelNode object at 0x7d05e2f9aad0>}
Channels:
{'__start__': <langgraph.channels.ephemeral_value.EphemeralValue object at 0x7d05e2c906c0>, '__end__': <langgraph.channels.last_value.LastValue object at 0x7d05e2c90c40>, '__previous__': <langgraph.channels.last_value.LastValue object at 0x7d05e1007280>}