以编程方式使用文档

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 是默认的通道类型。它存储写入的最新值,覆盖任何先前的值。将其用于输入和输出值,或用于在步骤之间传递数据。

from langgraph.channels import LastValue

channel: LastValue[int] = LastValue(int)

Topic(主题)

Topic 是一个可配置的 PubSub 通道,用于在参与者之间发送多个值或在步骤之间累积输出。它可以配置为对值进行去重,或累积在运行期间写入的所有值。

from langgraph.channels import Topic

# Accumulate all values written across steps
channel: Topic[str] = Topic(str, accumulate=True)

BinaryOperatorAggregate(二元运算符聚合)

BinaryOperatorAggregate 存储一个持久值,该值通过将二元运算符应用于当前值和每个新更新来进行更新。将其用于计算跨步骤的持续聚合。

from langgraph.channels import BinaryOperatorAggregate

# Running total: each write adds to the current value
total = BinaryOperatorAggregate(int, operator.add)

DeltaChannel(增量通道)

DeltaChannel 仅存储每个步骤中的增量,而不是完整的累积值。这对于频繁写入且随时间累积大量值的通道最有用——例如,长时间运行线程中的对话消息列表。如果没有增量存储,完整列表会在每个检查点重新序列化;而使用 DeltaChannel,则仅存储每个步骤中写入的新消息。

使用 DeltaChannel in an Annotated 类型注解的方式与使用普通 reducer 相同:

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 要求

传递给 reducerDeltaChannel is a **批量 reducer**:它接收当前状态和一个 *序列* 的所有写入(在单次调用中——不像标准 reducer 那样成对处理)。这与使用 Annotated in a StateGraph时使用的每键 reducer 不同,后者每次更新调用一次 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

两者都是可结合的:逐个应用批量与一次性应用批量会产生相同的结果。

使用快照_频率来限制读取延迟

没有快照的情况下,读取一个 DeltaChannel 值需要重放完整的写入历史——对于具有 N 步的线程为 O(N)。设置 snapshot_frequency=K 每 K 个 pregel 步写入一个完整的快照,将读取深度限制在最多 K 步:

class State(TypedDict):
    messages: Annotated[
        list[str],
        DeltaChannel(my_reducer, snapshot_frequency=5),
    ]

更高的 snapshot_frequency 值会减少存储开销但增加读取延迟。更低的值以更大的检查点为代价更严格地限制延迟。 None (默认值)完全跳过快照——适用于读取稀少或线程较短的情况。

版本兼容性和回滚

示例

虽然大多数用户会通过 StateGraph API 或 @entrypoint 装饰器与 Pregel 交互,但也可以直接与 Pregel 交互。

以下是几个不同的示例,让你了解 Pregel API 的用法。

Single node

    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"})
    
    {'b': 'foofoo'}
    

Multiple nodes

    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'}
    

Topic

    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

此示例演示如何使用 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' }
    

Cycle

此示例演示如何通过以下方式在图中引入循环 一条链写入其订阅的 channel。执行将继续进行 直到 None 值被写入 channel。

    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(图 API)函数式 API.

StateGraph (Graph API)

StateGraph(图 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 实例将与节点和通道列表关联。你可以通过打印来检查节点和通道。

    print(graph.nodes)
    

你将看到类似以下内容:

    {'__start__': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1810>,
     'write_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba14d0>,
     'score_essay': <langgraph.pregel.read.PregelNode at 0x7d05e3ba1710>}
    
    print(graph.channels)
    

你应该看到类似以下内容

    {'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

函数式 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>}