LangGraph 中的 StateGraph 本质上是:
节点读取 State
⬇
执行计算、调用LLM、调用工具
⬇
返回 Partial<State>
⬇
Reducer 决定如何把返回值合并回旧 State
每个 State 字段都可以配置自己的 Reducer, reducer 的形式是:
(old_value, new_value) -> merged_valuestate
为什么需要 State
LangGraph 中第一步需要定义一个 State, 为什么要这么做?
普通函数调用可以直接通过参数传递数据:
result_a = step_a(input)
result_b = step_b(result_a)
result_c = step_b(result_b)当 LangGraph 面临的问题复杂:
- 节点可能分支 / 循环
- 多个节点可能并行执行
- 节点可能中断并恢复
- 执行过程需要持久化
- 某些字段需要覆盖, 某些字段需要累积
因此, 不能只传递一个固定的 “上一步输出”.
LangGraph 引入一个共享的数据容器, 可以通过如下定义
- TypedDict: 一般情况够用
- dataclass: 可以设置默认值
- pydantic: 支持数据验证, 性能略差
所有节点都可以:
def node(state: State):
# 读取当前状态
...
# dosomething
# 返回希望更新的字段
return {
"some_field": new_value
}state 可以存什么
理论上, state 可以存储任何可序列化的数据, 例如:
class AgentState(TypedDict):
# LLM 对话上下文
messages: list
# 用户输入
query: str
# 工作流控制
retry_count: int
current_stage: str
# 检索结果
documents: list[dict]
# 结构化中间结果
extracted_entities: list[dict]
# 错误信息
errors: list[str]
# 最终输出
answer: str | None注意: state 应该存储 “需要跨节点共享、恢复或审计的数据”, 不是把所有局部变量都存进去
state 不等于持久化
容易混淆的点:
- state 是数据模型
- Checkpointer 是持久化机制
定义 state 并不意味着多次 invoke() 之间自动保留数据
只有配置 checkpointer, 并使用稳定的 thread_id, 图状态才能跨调用保存和恢复
from langgraph.checkpoint.memory import InMemorySaver
checkpointer = InMemorySaver()
graph = builder.compile(
checkpointer=checkpointer
)
config = {
"configurable": {
"thread_id": "user-123"
}
}
graph.invoke(
{"messages": [{"role": "user", "content": "你好"}]},
config=config,
)- State schema: 定义有哪些字段
- Reducer: 定义字段如何更新
- Checkpointer: 定义状态是否以及如何持久化
reducer
reducer 的作用
假设旧 State 是:
{
"documents": ["doc-1", "doc-2"]
}节点返回:
{
"documents": ["doc-2", "doc-3"]
}LangGraph 必须回答: 最终 documents 是什么?
可能有多种语义:
- 覆盖
- 追加
- 去重合并
- 只保留前10个
- 存在则修改, 不存在则追加
reducer 就是定义这个字段的状态更新规则
默认的 reducer
一个 state 字段, 若未声明 reducer, 默认为替换 (节点返回的新值会替换旧值)
from typing_extensions import TypedDict
class State(TypedDict):
count: int
results: list[str]节点:
def node_a(state: State):
return {
"count": 1,
"results": ["A"]
}
def node_b(state: State):
return {
"count": 2,
"results": ["B"]
}按顺序执行后的结果是:
{
"count": 2,
"results": ["B"]
}最简单的reducer: operator.add
reducer 通常通过 Annotated 声明:
from typing import Annotated
from typing_extensions import TypedDict
import operator
class State(TypedDict):
results: Annotated[list[str], operator.add]这表示: result = old_results + new_results, 即追加
自定义 reducer
reducer 是一个具有两个位置参数的函数:
- 左参数: 该字段的当前值
- 右参数: 节点返回的更新
当节点返回时, LangGraph 为每个更新的键调用 reducer, 并将返回值保存为新的状态值:
new_value = reducer(left=current_state[key], right=node_update[key])一个 reducer 应尽量满足下面几个性质:
1. 输入类型额输出类型一致
def reducer(old: list[str], new: list[str]) -> list[str]:
...2. 尽量不要原地修改旧值
# 不推荐
def reducer(old, new):
old.extend(new)
return old
# 更合适
def reducer(old, new):
return old + new原地修改会使状态快照、调试和重放更难分析
3.并行场景应尽量满足结合律
merge(merge(A, B), C) = merge(A, merge(B, C)并行分支的结果可能以不同分组的方式归并, 如果 reducer 不满足结合律, 结果可能依赖内部执行细节
4. 如果需要确定性, 最好满足交换律
# 理想清卡滚下:
merge(A, B) = merge(B, A)
# 数值加法满足:
1+2 = 2+1
# 列表拼接不满足:
["A"] + ["B"] != ["B"] + ["A"]这意味着并行节点使用列表追加 reducer 时, 不应依赖分支结果的绝对顺序
如果顺序很重要, 可以让数据包含显示排序字段, 最后在汇总节点中统一排序
{ "source_order": 1, "value": "A" }5. 注意幂等性
如果节点因为重试而产生相同更新:
new = ["Sigma-Aldrich"]
简单追加可能得到:
[ "Sigma-Aldrich", "Sigma-Aldrich" ]
如果业务要求幂等, 应使用去重 reducer, 或为对象设置唯一 ID
内置reducer: add_messages
为了保存大模型对话记录, 一般会在 state 定义一个 message 列表, 可以使用 operator.add 作为 reducer 来实现消息追加. 但对于一些特殊需求, 简单的追加就不够用了:
- 可能修改已有消息
- 可能删除指定消息 (人机中断)
- 可能传入字典格式
- 可能需要格式转换
- 同一个消息ID不应该重复追加
因此 LangGraph 提供了专用的 reducer, add_messages
from langgraph.graph.message import add_messages
from langchain.messages import AnyMessage
from typing import Annotated
from typing_extensions import TypedDict
class State(TypedDict):
messages: Annotated[list[AnyMessage], add_messages]add_messages 基本行为是:
- 新ID: 追加
- 相同ID: 更新原消息
- 支持消息格式规范化
- 配合 RemoveMessage 删除消息