从本文开始进入 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,最基础的组成部分有三个:
- State:保存节点之间需要共享的数据。
- Node:读取 State,执行一段逻辑,再返回状态更新。
- Edge:决定一个 Node 完成后应该执行哪个 Node。

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,例如 city、temperature 和 result 分别是三个 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]
city 和 temperature 是调用图时需要提供的字段;category 和 result 由节点逐步产生,因此使用 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 内部可以读写的全部字段。OutputState:invoke()最终返回给调用者的字段。
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 还可以接收 config 或 runtime,也可以使用 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:
- State Schema 定义工作流中可以共享哪些数据。
- Node 读取当前 State,并只返回自己负责的局部更新。
- Reducer 决定同一个 State 字段的新旧值如何合并。
- TypedDict、dataclass 和 Pydantic 适合不同的状态建模需求。
- 同步函数、异步函数和可调用对象都可以注册为 Node。
- 并行节点写入同一个字段时,需要为该字段配置合适的 Reducer。
这些内容解决了“工作流处理什么数据”和“节点怎样处理数据”两个问题。系列 10(下) 将继续介绍 Edge、条件路由、图的编译执行,并把这些知识组合成一个可以在 LangGraph Studio 中测试的城市旅行方案 Workflow。