LangGraph 系列 10(上):Workflow 的 State、Reducer 与 Node


从本文开始进入 LangGraph 系列的 Workflow 阶段。面对执行步骤明确、分支规则固定的业务,通常不需要让大模型决定每一步,而是由程序员把执行路径直接编排成 Workflow。

本篇是系列 10 的上篇,不调用大模型,而是使用普通 Python 函数依次学习 State、Reducer 和 Node。重点是理解数据如何在工作流中共享、更新和被节点处理,为下篇的 Edge、条件路由与完整案例打好基础。

本文对应代码位于:

llm-learning/langgraph/p10_workflow_basics/

1. Workflow 与 Agent 的区别

Workflow 和 Agent 都可以组织复杂任务,但二者决定执行路径的方式不同。

对比项 Workflow Agent
路径由谁决定 程序员通过 Node 和 Edge 提前定义 模型根据上下文动态决定
可预测性 高,同一输入通常经过相同路径 相对较低,模型可能采取不同动作
适用任务 审批、数据处理、规则判断、固定业务流程 开放式问答、动态工具选择、复杂推理
是否必须使用大模型 不需要 通常需要

例如,旅行方案 Workflow 可以固定为:

清理输入 → 校验输入 → 判断天气 → 并行生成方案 → 汇总结果

Agent 则可能先查询天气,再决定是否查询景点、酒店或交通信息。具体调用哪些工具、调用多少次,由模型在运行时决定。

Workflow 和 Agent 并不互斥。一个 Agent 可以成为 Workflow 中的某个 Node,多个 Agent 也可以由一张更大的 Workflow 统一编排。

2. Workflow 的组成

LangGraph 使用图来描述 Workflow,最基础的组成部分有三个:

  1. State:保存节点之间需要共享的数据。
  2. Node:读取 State,执行一段逻辑,再返回状态更新。
  3. Edge:决定一个 Node 完成后应该执行哪个 Node。

LangGraph Workflow 中 State、Node、Edge 与 Reducer 的关系

Reducer 不是第四种图元素。它是 State 中某个字段的更新规则,负责决定节点返回的新值应该覆盖、追加还是以其他方式合并到旧值中。

一张 StateGraph 从定义到执行,通常会经过下面的过程:

定义 State
→ 创建 StateGraph
→ 添加 Node
→ 添加 Edge
→ compile()
→ invoke() / stream()

接下来按照 State、Reducer、Node、Edge、路由函数的顺序逐一介绍。

3. State:工作流中的共享状态

State 表示一次 Workflow 运行过程中需要共享的数据。输入数据、中间结果、路由依据和最终结果都可以保存在 State 中。

3.1 State Schema、Channel 与状态快照

创建 StateGraph 时需要提供一个 State Schema:

# TravelState 决定图中允许使用哪些 Channel。
# 此时得到的是 Builder,还不能直接 invoke()。
builder = StateGraph(TravelState)

State Schema 描述整张图允许使用哪些字段,以及每个字段的数据类型。Schema 中的每个字段都可以理解为一个 Channel,例如 citytemperatureresult 分别是三个 Channel。

Node 不会直接把整个 State 替换掉,而是向一个或多个 Channel 提交局部更新。一个 Super-step 执行完成后,LangGraph 会把该轮更新合并成新的状态快照,再交给下一轮节点。

因此可以这样理解:

  • State Schema 定义状态的整体结构。
  • Channel 对应 State 中的单个字段。
  • State snapshot 是某个执行时刻各个 Channel 当前值的集合。

3.2 State 的三条基本规则

第一,Node 和路由函数可以读取当前 State:

def normalize_input(state: TravelState) -> dict:
    """读取当前 State,只返回自己负责的局部更新。"""

    return {"city": state["city"].strip()}

第二,Node 只需要返回本次要更新的字段,不需要返回完整 State:

return {"city": state["city"].strip()}

没有出现在返回值中的字段会保持原值。

第三,不要直接修改传入的 State。下面这种写法虽然像普通字典操作,但会让节点的输入和输出边界变得模糊:

# 不推荐
state["city"] = state["city"].strip()
return state

推荐始终返回一个新的局部更新:

# 推荐
return {"city": state["city"].strip()}

3.3 使用 TypedDict 定义 State

TypedDict 是最简单、最常用的 State 定义方式:

from typing_extensions import NotRequired, TypedDict

class TravelState(TypedDict):
    # 调用图时需要提供的输入字段。
    city: str
    temperature: int
    # 由 Node 在执行过程中逐步补充的字段。
    category: NotRequired[str]
    result: NotRequired[str]

citytemperature 是调用图时需要提供的字段;categoryresult 由节点逐步产生,因此使用 NotRequired

需要注意,TypedDict 主要用于类型提示。它不会在运行时自动检查温度是不是整数,也不会自动拒绝缺少字段的普通字典。真正的业务校验仍然需要由节点完成。

3.4 使用 dataclass 定义 State

dataclass 适合需要默认值、对象属性访问方式的 State:

from dataclasses import dataclass

@dataclass
class DataclassState:
    """dataclass 支持默认值,Node 通过属性访问字段。"""

    name: str
    greeting: str = ""

def greet(state: DataclassState) -> dict:
    # 使用 dataclass 时,Node 仍然返回字典形式的局部更新。
    return {"greeting": f"你好,{state.name}!"}

Node 收到的是 dataclass 实例,因此通过 state.name 读取字段。默认值可以被节点读取;如果希望某个默认字段稳定出现在最终结果中,最好由输入或节点显式写入该字段。

3.5 使用 Pydantic 定义 State

Pydantic 适合需要对初始输入进行运行时校验的场景:

from pydantic import BaseModel, Field

class PydanticState(BaseModel):
    name: str
    # 初始 State 创建时,age 必须在 0 到 120 之间。
    age: int = Field(ge=0, le=120)
    greeting: str = ""

如果初始输入中的 age-1,图会在创建初始 State 时抛出 ValidationError

这里有一个容易误解的边界:Pydantic 会校验初始输入,但 LangGraph 不会在每次 Node 更新后重新创建并校验完整 Pydantic 对象。业务规则仍然应该由节点明确校验,不能把 Pydantic 当作整个 Workflow 的持续校验器。

三种 Schema 可以这样选择:

Schema 优点 注意事项
TypedDict 简单、开销小、适合教学和大多数状态图 只提供类型提示,不做运行时校验
dataclass 支持默认值和属性访问 默认字段不一定自动出现在最终字典中
Pydantic 初始输入支持运行时校验 有额外开销,节点更新不会自动持续校验

3.6 分离输入、内部状态和输出

相对成熟的 Workflow 通常不希望调用者填写内部字段,也不希望把所有中间数据都暴露为最终输出。可以分别定义三种 Schema:

# 调用者允许传入的公开字段。
class InputState(TypedDict):
    text: str

# 图内部可以读写的完整 State。
class InternalState(TypedDict):
    text: str
    normalized_text: NotRequired[str]
    length: NotRequired[int]

# invoke() 执行完成后向调用者公开的字段。
class OutputState(TypedDict):
    normalized_text: str
    length: int

创建图时传入:

# 第一个参数是内部完整 State,另外两个参数限制公开输入和输出。
builder = StateGraph(
    InternalState,
    input_schema=InputState,
    output_schema=OutputState,
)

三者的职责是:

  • InputState:调用者允许传入的字段。
  • InternalState:整个 Workflow 内部可以读写的全部字段。
  • OutputStateinvoke() 最终返回给调用者的字段。

output_schema 只过滤最终输出。使用 stream_mode="values" 观察执行过程时,仍然可以看到内部状态快照。

完整示例位于 01_state_schemas.py,运行方式如下:

代码:

"""演示 TypedDict、dataclass、Pydantic 以及独立输入/输出 State。

State 是节点之间共享的数据,State Schema 则声明了这些数据包含哪些字段。
本文件先比较三种常见的 Schema 写法,再演示如何限制图的公开输入和输出。
"""

from dataclasses import dataclass

from langgraph.graph import END, START, StateGraph
from pydantic import BaseModel, Field, ValidationError
from typing_extensions import NotRequired, TypedDict


class TypedDictState(TypedDict):
    """TypedDict 适合定义轻量、清晰的 State。

    TypedDict 主要用于静态类型提示,不会在运行时自动校验数据。
    NotRequired 表示 greeting 可以不出现在图的初始输入中。
    """

    name: str
    greeting: NotRequired[str]


def greet_typed_dict(state: TypedDictState) -> dict:
    """Node 读取当前 State,只返回自己负责更新的 greeting 字段。"""

    return {"greeting": f'你好,{state["name"]}!'}


@dataclass
class DataclassState:
    """dataclass 可以提供默认值,节点通过属性访问 State 字段。"""

    name: str
    greeting: str = ""


def greet_dataclass(state: DataclassState) -> dict:
    """即使使用 dataclass,Node 仍然返回包含局部更新的字典。"""

    return {"greeting": f"你好,{state.name}!"}


class PydanticState(BaseModel):
    """Pydantic 会在运行时校验输入,适合对数据边界要求较高的场景。"""

    name: str
    age: int = Field(ge=0, le=120)
    greeting: str = ""


def greet_pydantic(state: PydanticState) -> dict:
    """Pydantic State 使用属性语法读取字段。"""

    return {"greeting": f"{state.name},你今年 {state.age} 岁。"}


# InputState 只描述调用者必须传入的字段。
class InputState(TypedDict):
    text: str


# InternalState 是图内部使用的完整 State,可以包含中间计算结果。
class InternalState(TypedDict):
    text: str
    normalized_text: NotRequired[str]
    length: NotRequired[int]


# OutputState 限制图最终暴露给调用者的字段。
class OutputState(TypedDict):
    normalized_text: str
    length: int


def normalize_text(state: InternalState) -> dict:
    """清理文本,并一次返回两个 State 字段的局部更新。"""

    normalized_text = state["text"].strip()
    return {
        "normalized_text": normalized_text,
        "length": len(normalized_text),
    }


def build_graph(state_schema, node):
    """创建只有一个 Node 的图,减少三个 Schema 示例的重复代码。

    StateGraph 是尚未编译的图构建器;START 和 END 是图的入口与出口。
    compile() 会检查图结构并返回真正可以 invoke() 的可执行图。
    """

    builder = StateGraph(state_schema)
    builder.add_node("greet", node)
    builder.add_edge(START, "greet")
    builder.add_edge("greet", END)
    return builder.compile()


typed_dict_graph = build_graph(TypedDictState, greet_typed_dict)
dataclass_graph = build_graph(DataclassState, greet_dataclass)
pydantic_graph = build_graph(PydanticState, greet_pydantic)

# 一个图可以分别声明内部、输入和输出 Schema。
schema_builder = StateGraph(
    InternalState,
    input_schema=InputState,
    output_schema=OutputState,
)
schema_builder.add_node("normalize_text", normalize_text)
schema_builder.add_edge(START, "normalize_text")
schema_builder.add_edge("normalize_text", END)
schema_graph = schema_builder.compile()


if __name__ == "__main__":
    # invoke() 传入初始 State,并等待整个图执行完成后返回最终 State。
    print("TypedDict:", typed_dict_graph.invoke({"name": "小明"}))
    print("dataclass:", dataclass_graph.invoke({"name": "小红"}))
    print("Pydantic:", pydantic_graph.invoke({"name": "小刚", "age": 20}))
    # text 属于内部 State,但不在 OutputState 中,因此不会出现在结果里。
    print("独立输入/输出 Schema:", schema_graph.invoke({"text": "  LangGraph  "}))

    # Pydantic 会在 Node 执行前拦截不满足 Field 约束的输入。
    try:
        pydantic_graph.invoke({"name": "小刚", "age": -1})
    except ValidationError as error:
        print("Pydantic 捕获非法年龄:", error.errors()[0]["msg"])

执行:

python langgraph/p10_workflow_basics/01_state_schemas.py

结果:

TypedDict: {'name': '小明', 'greeting': '你好,小明!'}
dataclass: {'name': '小红', 'greeting': '你好,小红!'}
Pydantic: {'name': '小刚', 'age': 20, 'greeting': '小刚,你今年 20 岁。'}
独立输入/输出 Schema: {'normalized_text': 'LangGraph', 'length': 9}
Pydantic 捕获非法年龄: Input should be greater than or equal to 0

4. Reducer:状态字段的操作规则

Node 返回的是局部状态更新,而 Reducer 决定这些更新如何写回对应 Channel。

Reducer 可以看作一个接收旧值和新值的二元函数:

new_value = reducer(current_value, update_value)

4.1 默认覆盖

没有显式声明 Reducer 时,字段采用默认覆盖规则:

class WorkflowState(TypedDict):
    # 没有通过 Annotated 绑定 Reducer,新值会覆盖旧值。
    score: int

第一个节点返回 {"score": 2},第二个节点返回 {"score": 3},最终值是 3,不是 5

默认覆盖适合分类结果、当前阶段、计算结果等只需要保留最新值的数据。

4.2 使用 operator.add 追加列表

如果多个节点需要向同一个列表追加内容,可以通过 Annotated 绑定 operator.add

import operator
from typing import Annotated

class WorkflowState(TypedDict):
    # Annotated 把 operator.add 绑定到 steps Channel。
    steps: Annotated[list[str], operator.add]

节点只返回本次新增的数据:

# 只返回本轮新增元素,不要把 State 中的旧列表一起返回。
return {"steps": ["normalize_input"]}

不要把旧列表和新内容一起返回,否则 Reducer 会再次追加旧内容,造成重复。

4.3 自定义 Reducer

当普通追加不能满足业务规则时,可以定义自己的 Reducer。例如合并标签、去重并排序:

def merge_unique(left: list[str], right: list[str]) -> list[str]:
    """合并旧值和本轮更新,返回去重、排序后的新列表。"""

    return sorted(set(left + right))

class WorkflowState(TypedDict):
    tags: Annotated[list[str], merge_unique]

Reducer 应尽量保持纯粹:相同输入产生相同输出,不修改传入的列表,也不执行网络请求或其他副作用。

4.4 使用 add_messages 合并消息

消息列表通常使用 LangGraph 提供的 add_messages

from langchain_core.messages import AnyMessage
from langgraph.graph.message import add_messages


class WorkflowState(TypedDict):
    # add_messages 会按消息 ID 合并,相同 ID 的消息会被替换。
    messages: Annotated[list[AnyMessage], add_messages]

add_messages 不只是普通追加。它会根据消息 ID 添加或替换消息,并能把支持的字典消息转换成 LangChain 消息对象。

消息增删改和 MessagesState 已经在 LangGraph 系列第 8 篇详细介绍,本篇只把它作为一种特殊 Reducer。

4.5 Reducer 的必须性

假设两个并行节点同时写入一个只有默认覆盖规则的字段:

class ConflictState(TypedDict):
    # 该字段没有 Reducer,并行写入时会发生冲突。
    value: str

两个节点分别返回 {"value": "A"}{"value": "B"}。LangGraph 无法判断应该保留哪个值,因此会抛出 InvalidUpdateError,而不是随机保留一个结果。

如果业务目标是合并多个列表,就必须显式声明:

# operator.add 允许同一 Super-step 中的多个列表更新被合并。
values: Annotated[list[str], operator.add]

4.6 Reducer 字段与最终输出字段分开

并行分支的完成顺序不应该被当作业务顺序保证。一个稳妥的设计是:

# 并行 Node 写入的临时 Channel,使用 Reducer 合并。
draft_sections: Annotated[list[PlanSection], operator.add]
# 汇总 Node 排序后写入的稳定输出,使用默认覆盖规则。
plan_sections: list[PlanSection]
  • 并行节点只向 draft_sections 追加带有 order 的原始区块。
  • 汇总节点对 draft_sections 排序。
  • 排序结果写入使用默认覆盖规则的 plan_sections

不要把排序结果重新写回 draft_sections,否则 operator.add 会把整份结果再次追加到原列表中。

完整示例位于 02_reducers.py

脚本会捕获一次故意制造的并行写入冲突,因此运行结束码仍然是 0

"""演示默认覆盖、列表追加、自定义 Reducer 和 add_messages。

Reducer 决定某个 State 字段收到更新时,旧值和新值应当如何合并。
它不是独立的图元素,而是通过 Annotated 配置在具体 State 字段上。
"""

import operator
import time
from typing import Annotated

from langchain_core.messages import AIMessage, AnyMessage, HumanMessage
from langgraph.errors import InvalidUpdateError
from langgraph.graph import END, START, StateGraph
from langgraph.graph.message import add_messages
from typing_extensions import TypedDict

def merge_unique(left: list[str], right: list[str]) -> list[str]:
    """自定义 Reducer:合并旧值 left 和新值 right,再去重、排序。"""

    return sorted(set(left + right))

class ReducerState(TypedDict):
    # 没有 Annotated 的字段使用默认 Reducer:新值直接覆盖旧值。
    score: int
    # operator.add 会把节点返回的列表追加到现有列表。
    steps: Annotated[list[str], operator.add]
    # Reducer 也可以是符合 (旧值, 新值) -> 合并值 形式的自定义函数。
    tags: Annotated[list[str], merge_unique]
    # add_messages 会按消息 ID 合并,支持追加消息以及替换同 ID 消息。
    messages: Annotated[list[AnyMessage], add_messages]


def first_update(state: ReducerState) -> dict:
    """第一个 Node 返回局部更新,不直接修改传入的 state。"""

    return {
        "score": 2,
        "steps": ["first_update"],
        "tags": ["python", "workflow"],
        "messages": [AIMessage(content="初始回答", id="answer-1")],
    }


def second_update(state: ReducerState) -> dict:
    """第二个 Node 再次更新相同字段,用于观察不同 Reducer 的行为。"""

    return {
        "score": 3,
        "steps": ["second_update"],
        "tags": ["langgraph", "python"],
        # 相同 ID 的消息会替换旧消息,而不是重复追加。
        "messages": [AIMessage(content="修正后的回答", id="answer-1")],
    }


# 固定 Edge 让两个 Node 顺序执行,便于对比旧值和新值如何合并。
builder = StateGraph(ReducerState)
builder.add_node(first_update)
builder.add_node(second_update)
builder.add_edge(START, "first_update")
builder.add_edge("first_update", "second_update")
builder.add_edge("second_update", END)
graph = builder.compile()


# 下面的图故意让两个并行 Node 更新同一个、没有 Reducer 的字段。
class ConflictState(TypedDict):
    value: str


def write_a(state: ConflictState) -> dict:
    print("    write_a 开始")
    time.sleep(0.2)
    print("    write_a 结束 -> value=A")
    return {"value": "A"}


def write_b(state: ConflictState) -> dict:
    print("    write_b 开始")
    time.sleep(0.2)
    print("    write_b 结束 -> value=B")
    return {"value": "B"}


conflict_builder = StateGraph(ConflictState)
conflict_builder.add_node(write_a)
conflict_builder.add_node(write_b)

# START 同时连接两个 Node,二者会在同一个 Super-step 中执行。
conflict_builder.add_edge(START, "write_a")
conflict_builder.add_edge(START, "write_b")

conflict_builder.add_edge("write_a", END)
conflict_builder.add_edge("write_b", END)

conflict_graph = conflict_builder.compile()


def demo_parallel_conflict() -> None:
    """用 stream 观察并行 Node 的执行,再在 State 合并阶段触发冲突。"""

    print("\n并行 Node 演示(START 同时连 write_a / write_b):")
    print("  若顺序执行,两次 sleep(0.2) 约需 0.4s;并发则约 0.2s。")

    started_at = time.perf_counter()
    node_updates: list[str] = []

    try:
        for update in conflict_graph.stream(
            {"value": "初始值"},
            stream_mode="updates",
        ):
            for node_name, node_output in update.items():
                elapsed = time.perf_counter() - started_at
                node_updates.append(node_name)
                print(
                    f"  [{elapsed:.2f}s] Super-step 收到 {node_name} 更新: "
                    f"{node_output}"
                )
    except InvalidUpdateError as error:
        elapsed = time.perf_counter() - started_at
        print(f"  [{elapsed:.2f}s] 两个 Node 都已执行,但合并 value 时发生冲突:")
        print(f"    {error}")
        print(f"  本步实际执行的 Node: {node_updates}")
        print("  说明:并行分支都跑完了,冲突发生在 Reducer 合并阶段,而不是某个 Node 没执行。")


if __name__ == "__main__":
    # 初始值也会参与 Reducer 合并,例如 steps 从 ["input"] 开始累加。
    result = graph.invoke( # 顺序执行,会先执行first_update,再执行second_update
        {
            "score": 1,
            "steps": ["input"],
            "tags": ["state"],
            "messages": [HumanMessage(content="请给我一条建议", id="user-1")],
        }
    )

    print("默认覆盖 score:", result["score"])
    print("operator.add steps:", result["steps"])
    print("自定义 Reducer tags:", result["tags"])
    print("messages:", result["messages"])
    print("add_messages:")
    for message in result["messages"]:
        print(" ", type(message).__name__, message.id, message.content)

    demo_parallel_conflict()

执行:

python langgraph/p10_workflow_basics/02_reducers.py

结果:

默认覆盖 score: 3
operator.add steps: ['input', 'first_update', 'second_update']
自定义 Reducer tags: ['langgraph', 'python', 'state', 'workflow']
messages: [HumanMessage(content='请给我一条建议', additional_kwargs={}, response_metadata={}, id='user-1'), AIMessage(content='修正后的回答', additional_kwargs={}, response_metadata={}, id='answer-1', tool_calls=[], invalid_tool_calls=[])]
add_messages:
  HumanMessage user-1 请给我一条建议
  AIMessage answer-1 修正后的回答

并行 Node 演示(START 同时连 write_a / write_b):
  若顺序执行,两次 sleep(0.2) 约需 0.4s;并发则约 0.2s。
    write_a 开始
    write_b 开始
    write_a 结束 -> value=A
  [0.20s] Super-step 收到 write_a 更新: {'value': 'A'}
    write_b 结束 -> value=B
  [0.20s] Super-step 收到 write_b 更新: {'value': 'B'}
  [0.20s] 两个 Node 都已执行,但合并 value 时发生冲突:
    At key 'value': Can receive only one value per step. Use an Annotated key to handle multiple values.
For troubleshooting, visit: https://docs.langchain.com/oss/python/langgraph/errors/INVALID_CONCURRENT_GRAPH_UPDATE
  本步实际执行的 Node: ['write_a', 'write_b']
  说明:并行分支都跑完了,冲突发生在 Reducer 合并阶段,而不是某个 Node 没执行。

5. Node:执行具体任务

Node 是 Workflow 中真正执行工作的步骤。它读取当前 State,完成一项职责,再返回局部状态更新。

5.1 同步函数 Node

普通同步函数适合纯计算或同步代码:

def normalize_text(state: NodeState) -> dict:
    """同步函数 Node:读取 State,返回局部更新。"""

    return {"normalized_text": state["text"].strip()}

如果 add_node() 只接收函数,LangGraph 默认使用函数名作为节点名:

# 未显式指定名称时,函数名 normalize_text 就是节点名。
builder.add_node(normalize_text)

也可以显式指定名称:

# 显式注册为 normalize,后续 Edge 应使用这个名称。
builder.add_node("normalize", normalize_text)

Edge 和路由函数使用的是注册后的节点名。

5.2 异步函数 Node

异步 Node 适合数据库、HTTP 请求等异步 I/O:

async def count_characters(state: NodeState) -> dict:
    """异步 Node 适合数据库、HTTP 请求等异步 I/O。"""

    length = len(state["normalized_text"])
    return {"length": length}

包含异步 Node 的图使用 ainvoke()astream()

# 图中包含异步 Node,因此使用 ainvoke()。
result = await graph.ainvoke({"text": "LangGraph"})

5.3 可调用对象 Node

实现了 __call__() 的对象也可以成为 Node:

class ResultFormatter:
    """可调用对象可以在初始化时保存可复用配置。"""

    def __init__(self, prefix: str):
        self.prefix = prefix

    def __call__(self, state: NodeState) -> dict:
        return {"result": f'{self.prefix}{state["normalized_text"]}'}

注册时可以把配置好的对象传给 add_node()

# 先创建带固定前缀的对象,再把它注册为 Node。
builder.add_node("format_result", ResultFormatter("处理结果"))

5.4 Node 的职责边界

一个 Node 最好只完成一项清晰的任务。例如校验输入、计算预算和整理报告应该拆成不同 Node,而不是堆在一个函数中。

根据需要,Node 还可以接收 configruntime,也可以使用 Runnable。本文只关注 StateGraph 的基础结构,不展开运行时上下文、Store 和工具调用。

完整示例位于 03_node_implementations.py

代码:

"""演示同步函数、异步函数和可调用对象三种 Node 写法。

Node 负责执行一小段业务逻辑:读取当前 State,并返回需要写回的局部更新。
多个职责清晰的小 Node,通常比一个包含全部逻辑的大函数更容易组合和测试。
"""

import asyncio

from langgraph.graph import END, START, StateGraph
from typing_extensions import NotRequired, TypedDict


class NodeState(TypedDict):
    # text 是输入字段,其余字段由后续 Node 逐步产生。
    text: str
    normalized_text: NotRequired[str]
    length: NotRequired[int]
    result: NotRequired[str]


def normalize_text(state: NodeState) -> dict:
    """同步函数 Node,适合无需等待外部 I/O 的计算。"""

    return {"normalized_text": state["text"].strip()}


async def count_characters(state: NodeState) -> dict:
    """异步函数 Node,适合数据库、HTTP 请求等异步 I/O。"""

    # sleep(0) 只用于模拟让出事件循环;真实项目中这里通常是 await I/O。
    await asyncio.sleep(0)
    return {"length": len(state["normalized_text"])}


class ResultFormatter:
    """实现 __call__ 的对象也可以作为 Node,并可保存初始化配置。"""

    def __init__(self, prefix: str):
        self.prefix = prefix

    def __call__(self, state: NodeState) -> dict:
        return {
            "result": (
                f'{self.prefix}{state["normalized_text"]}'
                f',共 {state["length"]} 个字符'
            )
        }


# add_node(函数) 默认使用函数名作为节点名,也可以显式传入节点名。
async_builder = StateGraph(NodeState)
async_builder.add_node(normalize_text)
async_builder.add_node("count_characters", count_characters)
async_builder.add_node("format_result", ResultFormatter("处理结果"))

# Edge 规定 Node 的固定执行顺序:入口 → 清理 → 计数 → 格式化 → 出口。
async_builder.add_edge(START, "normalize_text")
async_builder.add_edge("normalize_text", "count_characters")
async_builder.add_edge("count_characters", "format_result")
async_builder.add_edge("format_result", END)
async_graph = async_builder.compile()


# 再构建一个纯同步图,用于演示 invoke() 与 stream()。
class SyncState(TypedDict):
    value: int
    doubled: NotRequired[int]


def double_value(state: SyncState) -> dict:
    return {"doubled": state["value"] * 2}


sync_builder = StateGraph(SyncState)
sync_builder.add_node(double_value)
sync_builder.add_edge(START, "double_value")
sync_builder.add_edge("double_value", END)
sync_graph = sync_builder.compile()


async def main() -> None:
    # invoke() 适用于不包含异步 Node 的图,并一次返回最终 State。
    print("invoke 同步图:", sync_graph.invoke({"value": 6}))

    # updates 只输出每个 Node 本次写回的局部更新。
    print("stream updates:")
    for update in sync_graph.stream({"value": 6}, stream_mode="updates"):
        print(" ", update)

    # values 输出每一步合并后的完整 State 快照。
    print("stream values:")
    for value in sync_graph.stream({"value": 6}, stream_mode="values"):
        print(" ", value)

    # 图中含有 async Node 时,应使用 ainvoke(),避免阻塞事件循环。
    result = await async_graph.ainvoke({"text": "  LangGraph  "})
    print("ainvoke 异步图:", result)

    # astream() 是 stream() 的异步版本,使用 async for 消费事件。
    print("astream updates:")
    async for update in async_graph.astream(
        {"text": "  LangGraph  "},
        stream_mode="updates",
    ):
        print(" ", update)


if __name__ == "__main__":
    asyncio.run(main())

执行:

python langgraph/p10_workflow_basics/03_node_implementations.py

结果:

invoke 同步图: {'value': 6, 'doubled': 12}

stream updates:
  {'double_value': {'doubled': 12}}

stream values:
  {'value': 6}
  {'value': 6, 'doubled': 12}

ainvoke 异步图: {'text': '  LangGraph  ', 'normalized_text': 'LangGraph', 'length': 9, 'result': '处理结果:LangGraph,共 9 个字符'}

astream updates:
  {'normalize_text': {'normalized_text': 'LangGraph'}}
  {'count_characters': {'length': 9}}
  {'format_result': {'result': '处理结果:LangGraph,共 9 个字符'}}

5.5 分析 的执行

# 1. 图长什么样
START
normalize_text        (同步函数)
count_characters      (async 函数)
format_result         (ResultFormatter 实例)
END

初始 State 只有: {"text": "  LangGraph  "}

# 2. 逐步执行(4 个阶段)1 步:入口
LangGraph 收到输入,State 为:
{
    "text": "  LangGraph  "
}
按 Edge:START → normalize_text,调度第一个 Node。

第 2 步:normalize_text(同步)
def normalize_text(state):
    return {"normalized_text": state["text"].strip()}
读:state["text"]" LangGraph "
写回:{"normalized_text": "LangGraph"}
合并后 State:
{
    "text": "  LangGraph  ",
    "normalized_text": "LangGraph"
}
按 Edge:normalize_text → count_characters

第 3 步:count_characters(async)
async def count_characters(state):
    await asyncio.sleep(0)
    return {"length": len(state["normalized_text"])}
读:state["normalized_text"]"LangGraph"(9 个字符)
写回:{"length": 9}
合并后 State:
{
    "text": "  LangGraph  ",
    "normalized_text": "LangGraph",
    "length": 9
}
因为是 async Node,这里用 ainvoke 才能正确 await,不会阻塞事件循环。

按 Edge:count_characters → format_result

第 4 步:format_result(可调用对象)
注册时是:

add_node("format_result", ResultFormatter("处理结果"))
实际执行的是:

ResultFormatter("处理结果").__call__(state)
读:normalized_text、length
写回:
{"result": "处理结果:LangGraph,共 9 个字符"}
最终 State:

{
    "text": "  LangGraph  ",
    "normalized_text": "LangGraph",
    "length": 9,
    "result": "处理结果:LangGraph,共 9 个字符"
}
按 Edge:format_result → END,图结束。

第 5 步:返回结果
ainvoke 把最终 State 整包返回给 result,所以打印类似:

{
    "text": "  LangGraph  ",
    "normalized_text": "LangGraph",
    "length": 9,
    "result": "处理结果:LangGraph,共 9 个字符"
}

# 3. 几个容易混的点
Node 只返回「局部更新」
每个 Node 不会改整个 State,只返回自己负责的字段,LangGraph 负责合并:

return {"normalized_text": "..."}  # 不是 return 整个 state
State 在 Node 之间传递
Node A 写 normalized_text
    → 合并进 State
    → Node B 读到 normalized_text
    → 写 length
    → 合并进 State
    → Node C 读到 normalized_text + length
    → 写 result
为什么用 ainvoke 而不是 invoke
图里有 async Node(count_characters),应使用:

await async_graph.ainvoke(...)   # 正确
async_graph.invoke(...)          # 可能有问题或行为不一致
下面纯同步的 sync_graph 才适合用 invoke()# 4. 和 stream 的对比
同一套流程,如果改成:

async for update in async_graph.astream(..., stream_mode="updates"):
    print(update)
会分步看到每个 Node 的写回,而不是最后一次性拿完整 State:


{"normalize_text": {"normalized_text": "LangGraph"}}
{"count_characters": {"length": 9}}
{"format_result": {"result": "处理结果:LangGraph,共 9 个字符"}}

6. 总结

本篇从 Workflow 的基本组成出发,依次介绍了 State、Reducer 和 Node:

  1. State Schema 定义工作流中可以共享哪些数据。
  2. Node 读取当前 State,并只返回自己负责的局部更新。
  3. Reducer 决定同一个 State 字段的新旧值如何合并。
  4. TypedDict、dataclass 和 Pydantic 适合不同的状态建模需求。
  5. 同步函数、异步函数和可调用对象都可以注册为 Node。
  6. 并行节点写入同一个字段时,需要为该字段配置合适的 Reducer。

这些内容解决了“工作流处理什么数据”和“节点怎样处理数据”两个问题。系列 10(下) 将继续介绍 Edge、条件路由、图的编译执行,并把这些知识组合成一个可以在 LangGraph Studio 中测试的城市旅行方案 Workflow。


文章作者: hnbian
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 hnbian !
评论
  目录