外观
Pregel 实现了 LangGraph 的运行时,负责管理 LangGraph 应用程序的执行。
编译 StateGraph 或创建 @entrypoint 会产生一个可以用输入调用的 Pregel 实例。
Pregel 实现了 LangGraph 的运行时,负责管理 LangGraph 应用程序的执行。
编译 StateGraph 或创建 entrypoint 会产生一个可以用输入调用的 Pregel 实例。
本指南从较高层面解释运行时,并提供使用 Pregel 直接实现应用程序的说明。
注意:
Pregel运行时以 Google 的 Pregel 算法 命名,该算法描述了使用图进行大规模并行计算的高效方法。
注意:
Pregel运行时以 Google 的 Pregel 算法 命名,该算法描述了使用图进行大规模并行计算的高效方法。
概述
在 LangGraph 中,Pregel 将actor和**通道(channel)**组合到单个应用程序中。Actor 从通道中读取数据并向通道写入数据。Pregel 遵循 Pregel 算法/**整体同步并行(Bulk Synchronous Parallel)**模型,将应用程序的执行组织为多个步骤。
每个步骤包含三个阶段:
- 计划(Plan):确定在此步骤中执行哪些 actor。例如,在第一步中,选择订阅特殊输入通道的 actor;在后续步骤中,选择订阅上一步中更新过的通道的 actor。
- 执行(Execution):并行执行所有选中的 actor,直到全部完成、某个 actor 失败或达到超时。在此阶段,通道更新对 actor 不可见,直到下一步。
- 更新(Update):用本步骤中 actor 写入的值更新通道。
重复此过程,直到没有 actor 被选中执行,或达到最大步骤数。
Actor
actor 是一个 PregelNode。它订阅通道,从通道中读取数据,并向通道写入数据。可以将其视为 Pregel 算法中的 actor。PregelNode 实现了 LangChain 的 Runnable 接口。
通道(Channels)
通道用于在 actor(PregelNode)之间进行通信。每个通道都有一个值类型、一个更新类型和一个更新函数——该函数接收一组更新并修改存储的值。通道可用于将数据从一个链发送到另一个链,或在未来的步骤中将数据从一条链发送给它自身。
LastValue
LastValue 是默认的通道类型。它存储写入的最后一个值,覆盖之前的任何值。将其用于输入和输出值,或用于将数据从一个步骤传递到下一个步骤。
python
from langgraph.channels import LastValue
channel: LastValue[int] = LastValue(int)typescript
import { LastValue } from "@langchain/langgraph/channels";
const channel = new LastValue<number>();Topic
Topic 是一个可配置的 PubSub 通道,适用于在 actor 之间发送多个值,或跨步骤累积输出。它可以配置为对值进行去重,或累积一次运行期间写入的所有值。
python
from langgraph.channels import Topic
# 累积跨步骤写入的所有值
channel: Topic[str] = Topic(str, accumulate=True)typescript
import { Topic } from "@langchain/langgraph/channels";
// 累积跨步骤写入的所有值
const channel = new Topic<string>({ accumulate: true });BinaryOperatorAggregate
BinaryOperatorAggregate 存储一个持久化值,该值通过对当前值和每个新更新应用二元运算符进行更新。使用它来计算跨步骤的累计聚合值。
python
import operator
from langgraph.channels import BinaryOperatorAggregate
# 运行总计:每次写入都会累加到当前值
total = BinaryOperatorAggregate(int, operator.add)typescript
import { BinaryOperatorAggregate } from "@langchain/langgraph/channels";
// 运行总计:每次写入都会累加到当前值
const total = new BinaryOperatorAggregate<number>({ operator: (a, b) => a + b });DeltaChannel
WARNING
DeltaChannel 需要 langgraph>=1.2,目前处于测试版(beta)。API 在未来的版本中可能会发生变化。
DeltaChannel 在每个步骤中只存储增量差异,而不是完整的累积值。这对于频繁写入且随时间累积较大值的通道最为有用——例如,长时间运行线程中的对话消息列表。如果没有增量存储,完整列表会在每个检查点中被重新序列化;使用 DeltaChannel 时,只存储每个步骤写入的新消息。
TIP
当通道既被频繁写入又随时间不断变大时,请考虑使用 DeltaChannel。一个很好的信号:如果你注意到某个特定通道的检查点大小随线程长度线性增长,那么 DeltaChannel 很可能很适合。
在 Annotated 类型注解中使用 DeltaChannel 的方式与使用普通 reducer 相同:
python
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(bulk reducer):它单次调用即可接收当前状态和当前步骤所有写入的序列,而不是像标准 reducer 那样两两处理。这与 StateGraph 中与 Annotated 一起使用的按键 reducer 不同,后者的 reducer 在每次更新时被调用一次。
WARNING
批量 reducer 必须是可结合的(对批处理不变):
reducer(reducer(state, [xs]), [ys]) == reducer(state, [xs, ys])如果你的 reducer 不可结合,重建的状态可能会因 LangGraph 跨步骤批处理写入的方式而不同,从而产生不一致的行为。
WARNING
reducer 在重建时运行,而非写入时运行。 与 BinaryOperatorAggregate 不同(后者的 reducer 在写入时调用,因此合并后的值才是被序列化到检查点中的内容),DeltaChannel 的 reducer 是在通道值从其持久化的写入中重建时被调用的。被序列化的是原始的逐步写入;reducer 只在值被实体化时才被调用——即在下一次读取、下一个步骤的 actor,或重放历史时。
设计 reducer 时的实际影响:
- 使其成为
(state, writes)的纯函数。 任何副作用、随机性或墙钟时间读取(例如uuid.uuid4()、datetime.now())都会在值每次被重建时执行,并在每次重放时产生不同的结果。它们不会被固化到持久化的写入中。 - 不要依赖对传入写入的修改会被持久化。 如果你的 reducer 修改了写入对象(例如,为没有 ID 到达的条目分配一个稳定的 ID),该修改只存在于重建的值中。存储的写入仍然保持原始形态,因此下一次重建将再次看到未修改的输入。
- 在上游附加标识符和其他稳定元数据。 如果下游代码需要跨轮次按 ID 引用某个条目(例如,稍后更新或删除它),请在值写入通道之前分配该 ID——而不是在 reducer 内部。
以下是两种最常见情况的批量 reducer:
python
from typing import Any, Sequence
# List:按顺序追加所有写入
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:合并所有写入,键冲突时以最后一次写入为准
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 控制读取延迟
如果没有快照,读取 DeltaChannel 值需要重放完整的写入历史——对于具有 N 个步骤的线程为 O(N)。设置 snapshot_frequency=K 会每 K 个 pregel 步骤写入一次完整快照,将读取深度限制在至多 K 个步骤:
python
class State(TypedDict):
messages: Annotated[
list[str],
DeltaChannel(my_reducer, snapshot_frequency=5),
]snapshot_frequency 的值越大,存储开销越低,但读取延迟越高。值越小,延迟被约束得越紧,但代价是检查点更大。None(默认值)会完全跳过快照——适用于读取很少或线程较短的情况。
版本兼容性与回滚
WARNING
不支持回滚到不支持 DeltaChannel 的版本。 langgraph>=1.2 会以旧版本无法读取的新格式写入增量通道检查点。一旦某个线程使用了 DeltaChannel,降级 LangGraph 后这些检查点将无法读取,因为旧版运行时不理解增量格式,也无法重建通道状态。如果你需要回滚,请在降级之前使用 delta-channel-dump 恢复脚本迁移受影响的线程,或将其丢弃。
示例
虽然大多数用户会通过 StateGraph API 或 @entrypoint 装饰器与 Pregel 交互,但你也可以直接与 Pregel 交互。
虽然大多数用户会通过 StateGraph API 或 entrypoint 装饰器与 Pregel 交互,但你也可以直接与 Pregel 交互。
以下是几个不同的示例,让你对 Pregel API 有所了解。
单节点
python
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"})txt
{'b': 'foofoo'}typescript
import { EphemeralValue } from "@langchain/langgraph/channels";
import { Pregel, NodeBuilder } from "@langchain/langgraph/pregel";
const node1 = new NodeBuilder()
.subscribeOnly("a")
.do((x: string) => x + x)
.writeTo("b");
const app = new Pregel({
nodes: { node1 },
channels: {
a: new EphemeralValue<string>(),
b: new EphemeralValue<string>(),
},
inputChannels: ["a"],
outputChannels: ["b"],
});
await app.invoke({ a: "foo" });txt
{ b: 'foofoo' }多节点
python
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"})txt
{'b': 'foofoo', 'c': 'foofoofoofoo'}typescript
import { LastValue, EphemeralValue } from "@langchain/langgraph/channels";
import { Pregel, NodeBuilder } from "@langchain/langgraph/pregel";
const node1 = new NodeBuilder()
.subscribeOnly("a")
.do((x: string) => x + x)
.writeTo("b");
const node2 = new NodeBuilder()
.subscribeOnly("b")
.do((x: string) => x + x)
.writeTo("c");
const app = new Pregel({
nodes: { node1, node2 },
channels: {
a: new EphemeralValue<string>(),
b: new LastValue<string>(),
c: new EphemeralValue<string>(),
},
inputChannels: ["a"],
outputChannels: ["b", "c"],
});
await app.invoke({ a: "foo" });txt
{ b: 'foofoo', c: 'foofoofoofoo' }Topic
python
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"})txt
{'c': ['foofoo', 'foofoofoofoo']}typescript
import { EphemeralValue, Topic } from "@langchain/langgraph/channels";
import { Pregel, NodeBuilder } from "@langchain/langgraph/pregel";
const node1 = new NodeBuilder()
.subscribeOnly("a")
.do((x: string) => x + x)
.writeTo("b", "c");
const node2 = new NodeBuilder()
.subscribeTo("b")
.do((x: { b: string }) => x.b + x.b)
.writeTo("c");
const app = new Pregel({
nodes: { node1, node2 },
channels: {
a: new EphemeralValue<string>(),
b: new EphemeralValue<string>(),
c: new Topic<string>({ accumulate: true }),
},
inputChannels: ["a"],
outputChannels: ["c"],
});
await app.invoke({ a: "foo" });txt
{ c: ['foofoo', 'foofoofoofoo'] }BinaryOperatorAggregate
此示例演示如何使用 `BinaryOperatorAggregate` 通道来实现 reducer。
python
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"})txt
{ 'c': 'foofoo | foofoofoofoo' }typescript
import { EphemeralValue, BinaryOperatorAggregate } from "@langchain/langgraph/channels";
import { Pregel, NodeBuilder } from "@langchain/langgraph/pregel";
const node1 = new NodeBuilder()
.subscribeOnly("a")
.do((x: string) => x + x)
.writeTo("b", "c");
const node2 = new NodeBuilder()
.subscribeOnly("b")
.do((x: string) => x + x)
.writeTo("c");
const reducer = (current: string, update: string) => {
if (current) {
return current + " | " + update;
} else {
return update;
}
};
const app = new Pregel({
nodes: { node1, node2 },
channels: {
a: new EphemeralValue<string>(),
b: new EphemeralValue<string>(),
c: new BinaryOperatorAggregate<string>({ operator: reducer }),
},
inputChannels: ["a"],
outputChannels: ["c"],
});
await app.invoke({ a: "foo" });循环
此示例演示如何在图中引入循环,方法是让一条链写入它所订阅的通道。执行将继续,直到有 None 值写入该通道。
python
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"})txt
{'value': 'aaaaaaaaaaaaaaaa'}此示例演示如何在图中引入循环,方法是让一条链写入它所订阅的通道。执行将继续,直到有 null 值写入该通道。
typescript
import { EphemeralValue } from "@langchain/langgraph/channels";
import { Pregel, NodeBuilder, ChannelWriteEntry } from "@langchain/langgraph/pregel";
const exampleNode = new NodeBuilder()
.subscribeOnly("value")
.do((x: string) => x.length < 10 ? x + x : null)
.writeTo(new ChannelWriteEntry("value", { skipNone: true }));
const app = new Pregel({
nodes: { exampleNode },
channels: {
value: new EphemeralValue<string>(),
},
inputChannels: ["value"],
outputChannels: ["value"],
});
await app.invoke({ value: "a" });txt
{ value: 'aaaaaaaaaaaaaaaa' }高级 API
LangGraph 提供了两个用于创建 Pregel 应用程序的高级 API:StateGraph(Graph API) 和功能 API。
StateGraph(Graph API)
StateGraph(Graph API) 是一种更高层的抽象,可简化 Pregel 应用程序的创建。它允许你定义节点和边的图。当你编译图时,StateGraph API 会自动为你创建 Pregel 应用程序。
python
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")
# 编译图。
# 这将返回一个 Pregel 实例。
graph = builder.compile()StateGraph(Graph API) 是一种更高层的抽象,可简化 Pregel 应用程序的创建。它允许你定义节点和边的图。当你编译图时,StateGraph API 会自动为你创建 Pregel 应用程序。
typescript
import { START, StateGraph } from "@langchain/langgraph";
interface Essay {
topic: string;
content?: string;
score?: number;
}
const writeEssay = (essay: Essay) => {
return {
content: `Essay about ${essay.topic}`,
};
};
const scoreEssay = (essay: Essay) => {
return {
score: 10
};
};
const builder = new StateGraph<Essay>({
channels: {
topic: null,
content: null,
score: null,
}
})
.addNode("writeEssay", writeEssay)
.addNode("scoreEssay", scoreEssay)
.addEdge(START, "writeEssay")
.addEdge("writeEssay", "scoreEssay");
// 编译图。
// 这将返回一个 Pregel 实例。
const graph = builder.compile();编译后的 Pregel 实例将与一组节点和通道相关联。你可以通过打印它们来检查节点和通道。
python
print(graph.nodes)你将看到类似这样的内容:
txt
{'__start__': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1810>,
'write_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba14d0>,
'score_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1710>}python
print(graph.channels)你应该看到类似这样的内容:
txt
{'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>}typescript
console.log(graph.nodes);你将看到类似这样的内容:
txt
{
__start__: PregelNode { ... },
writeEssay: PregelNode { ... },
scoreEssay: PregelNode { ... }
}typescript
console.log(graph.channels);你应该看到类似这样的内容:
txt
{
topic: LastValue { ... },
content: LastValue { ... },
score: LastValue { ... },
__start__: EphemeralValue { ... },
writeEssay: EphemeralValue { ... },
scoreEssay: EphemeralValue { ... },
'branch:__start__:__self__:writeEssay': EphemeralValue { ... },
'branch:__start__:__self__:scoreEssay': EphemeralValue { ... },
'branch:writeEssay:__self__:writeEssay': EphemeralValue { ... },
'branch:writeEssay:__self__:scoreEssay': EphemeralValue { ... },
'branch:scoreEssay:__self__:writeEssay': EphemeralValue { ... },
'branch:scoreEssay:__self__:scoreEssay': EphemeralValue { ... },
'start:writeEssay': EphemeralValue { ... }
}功能 API
在功能 API 中,你可以使用 @entrypoint 来创建 Pregel 应用程序。entrypoint 装饰器允许你定义一个接收输入并返回输出的函数。
python
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)txt
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>}在功能 API 中,你可以使用 entrypoint 来创建 Pregel 应用程序。entrypoint 装饰器允许你定义一个接收输入并返回输出的函数。
typescript
import { MemorySaver } from "@langchain/langgraph";
import { entrypoint } from "@langchain/langgraph/func";
interface Essay {
topic: string;
content?: string;
score?: number;
}
const checkpointer = new MemorySaver();
const writeEssay = entrypoint(
{ checkpointer, name: "writeEssay" },
async (essay: Essay) => {
return {
content: `Essay about ${essay.topic}`,
};
}
);
console.log("Nodes: ");
console.log(writeEssay.nodes);
console.log("Channels: ");
console.log(writeEssay.channels);txt
Nodes:
{ writeEssay: PregelNode { ... } }
Channels:
{
__start__: EphemeralValue { ... },
__end__: LastValue { ... },
__previous__: LastValue { ... }
}