State、Node、Graph 三件套:状态模式、节点签名与追加 Reducer

本文是「LangGraph 核心与实战」系列专栏的第 2 篇。专栏总览参见:《LangGraph 核心架构全景:用状态图构建可控 Agent 工作流》。

任何 LangGraph 应用在底层都由三个核心部件构成:状态(State)、节点(Node)与图(Graph)。

虽然很多入门示例只有十来行代码,但在深入开发复杂系统时,开发者往往会遭遇几类典型困惑:为什么节点返回的字典不能包含原有的全部数据?为什么历史消息列表在第二步直接被覆盖成了最后一条?节点内部能不能访问外部全局变量?

理清这三者的契约规范与运行机制,是编写稳健状态图工作流的基础。


环境安装与基准依赖

LangGraph 是一个轻量级编排层,核心运行机制仅依赖 Python 标准库与基础类型支持:

1
2
# 核心依赖安装
pip install langgraph langchain-core

1. State(状态):整张图的共享内存

State 是图执行过程中的全局数据上下文。图中的每一个节点在被调用时,都会接收到当前的 State 实例;执行完毕后,节点输出的数据会再次合并回 State 中。

用 TypedDict 进行强类型数据建模

LangGraph 最推荐使用 Python 标准库中的 typing.TypedDict 来定义 State:

1
2
3
4
5
6
7
from typing import TypedDict, List, Optional

class ProcessingState(TypedDict):
raw_input: str # 原始用户输入
sanitized_text: str # 清洗后的文本
token_count: int # 计算得出的 Token/词数
summary: Optional[str] # 阶段总结,初始可为空

TypedDict 在代码编写期提供了 IDE 的自动补全与类型检查能力,同时在运行时保持了原生 Python 字典的低开销特性。

状态更新机制:局部增量合并(Delta Merge)

初学者最容易犯的错误是在节点函数中试图修改整个 State 并原样返回:

1
2
3
4
# 错误示范:试图返回全量 State
def bad_node(state: ProcessingState) -> ProcessingState:
state["sanitized_text"] = state["raw_input"].strip()
return state # 严禁原地修改并返回完整 state

LangGraph 的核心法则是:节点必须只返回“当前需要更新的字段增量字典”:

1
2
3
4
# 正确示范:只返回发生变化的键值
def good_node(state: ProcessingState) -> dict:
cleaned = state["raw_input"].strip()
return {"sanitized_text": cleaned}

LangGraph 在底层执行的是浅层字典合并(Shallow Merge):

  • 节点返回 {"sanitized_text": "..."} 时,系统仅覆盖 State 中的该键;
  • State 中的其他字段(如 raw_input、token_count)保持原样不变;
  • 这种机制避免了多节点传递中由于全量深拷贝引发的内存暴涨,也保证了数据变更记录的纯粹性。

列表字段的覆盖与追加:使用 Reducer

默认情况下,如果一个字段的值是列表,后续节点返回同名列表会直接覆盖前序节点的内容:

1
2
3
# 默认行为是覆盖 (Overwrite)
state["logs"] = ["step 1 ok"]
# 节点返回 {"logs": ["step 2 ok"]} 后,state["logs"] 变为 ["step 2 ok"],step 1 丢失!

为了让列表字段支持“追加(Append)”,必须使用 Python 的 Annotated 语法为该字段注入 归约函数(Reducer):

1
2
3
4
5
6
7
from typing import Annotated, List, TypedDict
import operator

class ConversationState(TypedDict):
# 使用 operator.add 声明该字段采用追加合并策略
history: Annotated[List[str], operator.add]
current_topic: str

当为字段附加了 operator.add 之后:

  • 节点 A 返回 {"history": ["用户: 你好"]};
  • 节点 B 返回 {"history": ["助手: 您好,有什么可以帮您?"]};
  • LangGraph 底层会自动执行 history = operator.add(history, new_items),状态最终聚合为包含两条记录的完整列表。

2. Node(节点):图的操作执行单元

节点是承载具体计算的逻辑实体。在代码层面,一个 Node 就是一个普通的 Python 函数。

节点签名的三大工程准则

标准节点函数的类型签名必须满足:

编写节点时应严格恪守以下三项原则:

  1. 原则一:单一职责(Single Responsibility):不要把“调外部 API”、“数据清洗”、“计算打分”与“持久化写库”写进同一个节点。节点越原子化,图在后续重构或插入重试边时就越灵活;
  2. 原则二:输入确定性(Pure Function):节点的内部逻辑应当尽量只读取 state 参数,避免在节点内读取外部未受控的全局变量。依赖外部服务时(如调用大模型或数据库),应通过外部依赖注入,便于单独进行单元测试;
  3. 原则三:局部输出:只向外输出确需写入全局上下文的数据,临时计算变量不要注入状态,防止状态体积过度膨胀。

3. Graph(图):装配拓扑与编译执行

有了状态契约与处理节点后,使用 StateGraph 将它们连接为计算拓扑:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
from langgraph.graph import StateGraph, START, END

# 1. 实例化图构造器,并绑定 State 类型
builder = StateGraph(ProcessingState)

# 2. 注册节点 (内部标识符, 函数引用)
builder.add_node("clean_text", clean_text_node)
builder.add_node("count_tokens", count_tokens_node)

# 3. 编排转移边 (Edge)
# 使用内置常量 START 作为图的唯一入口
builder.add_edge(START, "clean_text")
builder.add_edge("clean_text", "count_tokens")
# 使用内置常量 END 标记终止
builder.add_edge("count_tokens", END)

# 4. 编译图为可执行应用 (CompiledGraph)
app = builder.compile()

两种执行模式:invoke 与 stream

编译生成的 app 提供了两种运行方式:

模式 A:app.invoke()(批量同步交付)

适用于传统 Web API 接口或单次批处理,等待整个图跑完全部节点后,一次性返回最终收敛的完整 State:

1
2
final_state = app.invoke({"raw_input": "  这是一段需要清洗的测试数据   "})
print("最终结果:", final_state)

模式 B:app.stream()(流式逐步观测)

在调试排障、长程任务或需要向前端提供实时进度条的场景下,使用 stream 逐节点消费执行事件。每次 yield 产出一个字典,键为当前刚执行完的节点名称,值为该节点返回的增量更新:

1
2
3
for event in app.stream({"raw_input": "  这是一段需要清洗的测试数据   "}):
for node_name, state_update in event.items():
print(f"工位 [{node_name}] 完成,输出增量: {state_update}")

完整可运行实战:文本结构化处理流水线

以下代码整合了强类型 State、operator.add 增量日志追加、多节点顺序编排与 stream 观测:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
from typing import TypedDict, Annotated, List, Optional
import operator
from langgraph.graph import StateGraph, START, END

# 1. 定义数据契约
class PipelineState(TypedDict):
raw_content: str
cleaned_content: str
word_count: int
extracted_keywords: List[str]
# 执行审计轨迹日志:使用 operator.add 追加
audit_trace: Annotated[List[str], operator.add]

# 2. 编写原子工序节点
def sanitize_node(state: PipelineState) -> dict:
raw = state["raw_content"]
cleaned = " ".join(raw.strip().split())
return {
"cleaned_content": cleaned,
"audit_trace": ["sanitize_node: 已剔除多余空格与不可见字符"]
}

def analyze_node(state: PipelineState) -> dict:
cleaned = state["cleaned_content"]
words = cleaned.split(" ")
# 模拟提取超过 3 个字符的关键词
keywords = [w for w in words if len(w) > 3]
return {
"word_count": len(words),
"extracted_keywords": keywords,
"audit_trace": [f"analyze_node: 统计完成,词数={len(words)},关键词数={len(keywords)}"]
}

# 3. 编排图结构
builder = StateGraph(PipelineState)

builder.add_node("sanitize", sanitize_node)
builder.add_node("analyze", analyze_node)

builder.add_edge(START, "sanitize")
builder.add_edge("sanitize", "analyze")
builder.add_edge("analyze", END)

app = builder.compile()

# 4. 执行并观察审计轨迹
if __name__ == "__main__":
initial_input = {
"raw_content": " LangGraph provides stateful orchestration for LLM applications "
}

print("--- 启动流式执行 ---")
for event in app.stream(initial_input):
for node_id, delta in event.items():
print(f"[{node_id}] -> 产出字段: {list(delta.keys())}")

print("\n--- 获取最终收敛状态 ---")
final_output = app.invoke(initial_input)
print("最终清洗文本:", final_output["cleaned_content"])
print("关键词清单:", final_output["extracted_keywords"])
print("全局审计日志:")
for log in final_output["audit_trace"]:
print(f" * {log}")

运行输出:

1
2
3
4
5
6
7
8
9
10
--- 启动流式执行 ---
[sanitize] -> 产出字段: ['cleaned_content', 'audit_trace']
[analyze] -> 产出字段: ['word_count', 'extracted_keywords', 'audit_trace']

--- 获取最终收敛状态 ---
最终清洗文本: LangGraph provides stateful orchestration for LLM applications
关键词清单: ['LangGraph', 'provides', 'stateful', 'orchestration', 'applications']
全局审计日志:
* sanitize_node: 已剔除多余空格与不可见字符
* analyze_node: 统计完成,词数=7,关键词数=5

系列导航与参考