Pregel 实现了 LangGraph 的运行时,管理 LangGraph 应用程序的执行。
编译 StateGraph 或创建 entrypoint 会生成一个 Pregel 实例,可以接受输入进行调用。
本指南从高层介绍运行时,并提供直接使用 Pregel 实现应用程序的说明。
> **Note:** Pregel 运行时的命名源自 Google 的 Pregel 算法,该算法描述了一种使用图进行大规模并行计算的高效方法。
概述
在 LangGraph 中,Pregel 将 **参与者** 和 **通道** 组合成单个应用程序。 **参与者** 从通道读取数据并向通道写入数据。Pregel 按照 **Pregel 算法**/**批量同步并行** model.
模型组织应用程序的执行,每个步骤包含三个阶段:
- * **计划**:确定在此步骤中要执行的 **参与者** 。例如,在第一步中,选择订阅特殊 **输入** 通道的 **参与者** ;在后续步骤中,选择订阅上一步更新的通道的 **参与者** 。
- * **执行**:并行执行所有选中的 **参与者** ,直到全部完成,或其中一个失败,或达到超时。在此阶段,通道更新对参与者不可见,直到下一步。
- * **更新**:使用该步骤中 **参与者** 写入的值更新通道。
重复直到没有 **参与者** 被选中执行,或达到最大步骤数。
参与者
An **参与者** is a PregelNode它订阅通道,从通道中读取数据,并向通道写入数据。可以将其视为一个 **参与者(actor)** 在 Pregel 算法中。 PregelNodes 实现 LangChain 的 Runnable 接口。
通道(Channels)
通道用于在参与者(PregelNodes)之间通信。每个通道都有一个值类型、一个更新类型和一个更新函数——更新函数接收一系列更新并修改存储的值。通道可用于将数据从一个链发送到另一个链,或将数据从链发送到自身的后续步骤中。
LastValue(最新值)
LastValue 是默认的通道类型。它存储写入的最新值,覆盖任何先前的值。将其用于输入和输出值,或用于在步骤之间传递数据。
const channel = new LastValue<number>();
Topic(主题)
Topic 是一个可配置的 PubSub 通道,用于在参与者之间发送多个值或在步骤之间累积输出。它可以配置为对值进行去重,或累积在运行期间写入的所有值。
// Accumulate all values written across steps
const channel = new Topic<string>({ accumulate: true });
BinaryOperatorAggregate(二元运算符聚合)
BinaryOperatorAggregate 存储一个持久值,该值通过将二元运算符应用于当前值和每个新更新来进行更新。将其用于计算跨步骤的持续聚合。
// Running total: each write adds to the current value
const total = new BinaryOperatorAggregate<number>({ operator: (a, b) => a + b });
示例
虽然大多数用户会通过 StateGraph API 或 entrypoint 装饰器与 Pregel 交互,但也可以直接与 Pregel 交互。
以下是几个不同的示例,让你了解 Pregel API 的用法。
Single node
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" });
{ b: 'foofoo' }
Multiple nodes
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" });
{ b: 'foofoo', c: 'foofoofoofoo' }
Topic
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" });
{ c: ['foofoo', 'foofoofoofoo'] }
BinaryOperatorAggregate
此示例演示如何使用 BinaryOperatorAggregate channel 实现一个 reducer。
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" });
Cycle
本示例演示了如何通过以下方式在图中引入循环 一条链向其订阅的通道写入数据。执行将继续进行 直到 null 值被写入通道为止。
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" });
{ value: 'aaaaaaaaaaaaaaaa' }
高级 API
LangGraph 提供了两个用于创建 Pregel 应用的高级 API: StateGraph(图 API) 和 函数式 API.
StateGraph (Graph API)
StateGraph(图 API) 是一个更高级的抽象,简化了 Pregel 应用的创建。它允许你定义节点和边的图。当编译图时,StateGraph API 会自动为你创建 Pregel 应用。
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({
channels: {
topic: null,
content: null,
score: null,
}
})
.addNode("writeEssay", writeEssay)
.addNode("scoreEssay", scoreEssay)
.addEdge(START, "writeEssay")
.addEdge("writeEssay", "scoreEssay");
// Compile the graph.
// This will return a Pregel instance.
const graph = builder.compile();
编译后的 Pregel 实例将与节点和通道列表关联。你可以通过打印来检查节点和通道。
console.log(graph.nodes);
你将看到类似以下内容:
{
__start__: PregelNode { ... },
writeEssay: PregelNode { ... },
scoreEssay: PregelNode { ... }
}
console.log(graph.channels);
你应该看到类似以下内容
{
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 { ... }
}
Functional API
在 函数式 API中,你可以使用 entrypoint 来创建 Pregel 应用。 entrypoint 装饰器允许你定义一个接受输入并返回输出的函数。
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);
Nodes:
{ writeEssay: PregelNode { ... } }
Channels:
{
__start__: EphemeralValue { ... },
__end__: LastValue { ... },
__previous__: LastValue { ... }
}