本文档系统阐述 LangGraph 框架的核心概念与实战模式,涵盖从图基础、状态管理、条件路由到持久化恢复、人机交互、子图并行、流式输出、多 Agent 编排和生产级 Agent 的完整技术链路,适合作为使用 LangGraph 构建复杂 Agent 工作流的核心技术参考资料。
1. 图基础
1.1 定义
LangGraph 是基于状态图(StateGraph)的 Agent 编排框架。State(状态)、Node(节点)、Edge(边)和 Conditional Edge(条件路由)是构建图的四要素。StateGraph 是有向图,节点是函数,边定义执行顺序,条件边实现分支逻辑--这是 LangGraph 一切的基础。
1.2 核心概念
+------------------------------------------------------------------+
| LangGraph 核心概念 |
+------------------------------------------------------------------+
| |
| StateGraph (状态图) |
| | |
| +-- State (状态): 在节点间传递的数据容器 |
| | - TypedDict 定义结构 |
| | - Reducer 控制合并方式 |
| | |
| +-- Node (节点): 处理状态的函数 |
| | - 接收 State,返回部分 State 更新 |
| | - 每个节点是一个 Python 函数 |
| | |
| +-- Edge (边): 定义节点间的执行顺序 |
| | - 普通边: A -> B (固定顺序) |
| | - 条件边: A -> (条件判断) -> B or C |
| | |
| +-- START / END: 特殊节点 |
| - START: 图的入口 |
| - END: 图的出口 |
| |
+------------------------------------------------------------------+1.3 基础图构建
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# Step 1: 定义 State
class BasicState(TypedDict):
"""图的状态:在节点间传递的数据容器"""
user_input: str
messages: Annotated[list[str], lambda x, y: x + y] # 使用 reducer 累加
current_step: str
result: str
# Step 2: 定义节点函数
def receive_input(state: BasicState) -> dict:
"""节点:接收用户输入"""
return {"messages": [f"[系统] 收到输入: {state['user_input']}"],
"current_step": "received"}
def process_data(state: BasicState) -> dict:
"""节点:处理数据"""
processed = f"已处理: {state['user_input']}"
return {"messages": [f"[处理] {processed}"],
"current_step": "processed",
"result": processed}
def generate_output(state: BasicState) -> dict:
"""节点:生成输出"""
return {"messages": [f"[输出] 最终结果: {state.get('result', '无')}"],
"current_step": "done"}
# Step 3: 构建图
def build_basic_graph():
"""构建基础线性图: START -> receive -> process -> generate -> END"""
g = StateGraph(BasicState)
g.add_node("receive_input", receive_input)
g.add_node("process_data", process_data)
g.add_node("generate_output", generate_output)
g.add_edge(START, "receive_input")
g.add_edge("receive_input", "process_data")
g.add_edge("process_data", "generate_output")
g.add_edge("generate_output", END)
return g.compile()1.4 图的执行流程
图执行流程:
invoke({"user_input": "hello"})
|
v
START
|
v
+------------------+
| receive_input | state["messages"] += ["收到输入"]
| (节点函数) | state["current_step"] = "received"
+--------+---------+
|
v
+------------------+
| process_data | state["result"] = "已处理: hello"
| (节点函数) | state["messages"] += ["处理完成"]
+--------+---------+
|
v
+------------------+
| generate_output | state["messages"] += ["输出结果"]
| (节点函数) | state["current_step"] = "done"
+--------+---------+
|
v
END
|
v
最终 State (包含所有更新)1.5 节点函数的返回值规则
节点返回值规则:
节点函数返回 dict,只包含需要更新的字段:
def my_node(state) -> dict:
return {"field_a": "new_value"} # 只更新 field_a
规则:
- 返回的字段会根据 Reducer 合并到 State
- 未返回的字段保持不变
- 默认 Reducer = "覆盖" (overwrite)
- 带 Annotated 的字段使用指定 Reducer
示例:
State: {messages: ["msg1"], result: ""}
Node 返回: {messages: ["msg2"], result: "done"}
合并后 State: {messages: ["msg1", "msg2"], result: "done"}
^-- reducer 累加 ^-- 覆盖1.6 与其他概念的关联
-> 状态管理:State 的合并方式由 Reducer 控制
-> 条件路由与分支:在基础图上添加条件边
<- Agent 核心能力:LangGraph 是 Agent 核心能力的框架实现
2. 状态管理
2.1 定义
State 管理是 LangGraph 的核心。State Schema 定义数据结构,Reducer 控制状态合并方式。默认是"覆盖"模式,用 Annotated + reducer 可以实现"累加"模式。Messages State 是最常用的 reducer 模式,支持消息列表的自动累加。
对应 Demo: demos/04_LangGraph/02_状态管理.py
2.2 三种 State 合并模式
2.2.1 覆盖模式(默认)
class OverwriteState(TypedDict):
"""默认覆盖模式:节点返回的字段会覆盖旧值"""
current_value: str
history: list[str] # 注意:默认覆盖,不是累加!
def step_a(state: OverwriteState) -> dict:
return {"current_value": "A", "history": ["执行了A"]} # history 被覆盖!
def step_b(state: OverwriteState) -> dict:
return {"current_value": "B", "history": ["执行了B"]} # history 被覆盖!
# 执行 step_a -> step_b 后:
# current_value = "B" (覆盖)
# history = ["执行了B"] (覆盖,丢失了 ["执行了A"])2.2.2 累加模式(Reducer)
class AccumulateState(TypedDict):
"""累加模式:使用 reducer 实现状态累加"""
current_value: str # 覆盖模式
history: Annotated[list[str], lambda old, new: old + new] # 累加模式
counter: Annotated[int, lambda old, new: old + new] # 数字累加
def step_c(state: AccumulateState) -> dict:
return {"current_value": "C", "history": ["执行了C"], "counter": 1}
def step_d(state: AccumulateState) -> dict:
return {"current_value": "D", "history": ["执行了D"], "counter": 10}
# 执行 step_c -> step_d 后:
# current_value = "D" (覆盖)
# history = ["执行了C", "执行了D"] (累加!)
# counter = 11 (累加!)2.2.3 Messages State(消息列表)
from langgraph.graph.message import add_messages
from langchain_core.messages import AIMessage, HumanMessage
class MessagesState(TypedDict):
"""消息列表状态:使用 add_messages reducer"""
messages: Annotated[list, add_messages] # 自动处理消息累加和ID去重
def chat_node(state: MessagesState) -> dict:
"""对话节点"""
return {"messages": [AIMessage(content="你好!有什么可以帮你?")]}2.3 Reducer 工作原理
Reducer 合并示意图:
覆盖 Reducer (默认):
old_state[field] = new_value
-> 旧值被完全替换
累加 Reducer (list):
old_state[field] = old_state[field] + new_value
-> 新值追加到列表
add_messages Reducer:
-> 按 message ID 去重后追加
-> 相同 ID 的消息会被替换(更新)
自定义 Reducer:
def my_reducer(old, new):
return merge_logic(old, new)
-> 完全自定义合并逻辑2.4 自定义 Reducer
def merge_dicts(old: dict, new: dict) -> dict:
"""自定义 reducer: 字典深度合并"""
result = old.copy()
for key, value in new.items():
if key in result and isinstance(result[key], dict) and isinstance(value, dict):
result[key] = merge_dicts(result[key], value)
else:
result[key] = value
return result
class CustomState(TypedDict):
data: Annotated[dict, merge_dicts] # 使用自定义字典合并 reducer
logs: Annotated[list[str], lambda x, y: x + y]2.5 Reducer 模式对比
2.6 状态设计最佳实践
State 设计原则:
1. 最小化 State:
-> 只存必要信息,避免冗余
-> 大文本用引用/ID,不存全文
2. 合理使用 Reducer:
-> 需要累加的用 Annotated + reducer
-> 只需当前值的用默认覆盖
3. 分层 State:
-> 基础状态: messages, iteration
-> 任务状态: goal, plan, result
-> 控制状态: status, error, budget
4. 类型安全:
-> 使用 TypedDict 明确字段类型
-> 必填字段和可选字段区分2.7 与其他概念的关联
<- 图基础:State 是图的基础组件
-> 条件路由:条件边基于 State 字段做路由判断
-> ReAct Agent:Agent 状态用 Messages State 管理
3. 条件路由与分支
3.1 定义
条件路由是 LangGraph 实现分支逻辑的核心。add_conditional_edges 根据状态动态选择下一个节点,支持多路分支、循环控制、动态路径和基于 LLM 决策的路由,是 Agent 工作流的基础。
对应 Demo: demos/04_LangGraph/03_条件路由与分支.py
3.2 条件路由架构
+------------------------------------------------------------------+
| 条件路由架构 |
+------------------------------------------------------------------+
| |
| 基础多路分支: |
| |
| +---> handle_search (intent=search) |
| | |
| classify +---> handle_calculate (intent=calculate) |
| intent | |
| +---> handle_translate (intent=translate) |
| | |
| +---> clarify (confidence < 0.5) |
| |
| 循环重试: |
| |
| +-------+ +-------+ |
| | validate |--->| handle | |
| +---+---+ +---+---+ |
| | | |
| | (失败) | (成功) |
| +<--- retry <---+ |
| | |
| v |
| END |
| |
+------------------------------------------------------------------+3.3 多路分支实现
class RouterState(TypedDict):
"""路由状态"""
user_input: str
messages: Annotated[list[str], lambda x, y: x + y]
intent: str
confidence: float
retry_count: int
result: str
def classify_intent(state: RouterState) -> dict:
"""意图分类节点"""
user_input = state["user_input"]
if "搜索" in user_input or "查询" in user_input:
intent, confidence = "search", 0.9
elif "计算" in user_input or "算" in user_input:
intent, confidence = "calculate", 0.85
elif "翻译" in user_input:
intent, confidence = "translate", 0.8
else:
intent, confidence = "unknown", 0.3
return {"intent": intent, "confidence": confidence}
def route_by_intent(state: RouterState) -> str:
"""根据意图路由(返回节点名)"""
if state["confidence"] < 0.5:
return "clarify"
return state["intent"]
# 构建带条件路由的图
g = StateGraph(RouterState)
g.add_node("classify", classify_intent)
g.add_node("search", handle_search)
g.add_node("calculate", handle_calculate)
g.add_node("translate", handle_translate)
g.add_node("clarify", handle_clarify)
g.add_edge(START, "classify")
g.add_conditional_edges("classify", route_by_intent, {
"search": "search",
"calculate": "calculate",
"translate": "translate",
"clarify": "clarify",
})
g.add_edge("search", END)
g.add_edge("calculate", END)
g.add_edge("translate", END)
g.add_edge("clarify", END)3.4 循环重试实现
def validate_result(state: RouterState) -> dict:
"""验证结果"""
if not state.get("result"):
return {"retry_count": state.get("retry_count", 0) + 1}
def should_retry(state: RouterState) -> str:
"""决定是否重试"""
if state.get("result"):
return END
if state.get("retry_count", 0) >= 3:
return END
return "handle" # 重试
# 循环重试图
g.add_edge(START, "handle")
g.add_edge("handle", "validate")
g.add_conditional_edges("validate", should_retry, {
END: END,
"handle": "handle", # 循环回 handle
})3.5 条件路由 vs 固定边
3.6 路由函数设计
路由函数最佳实践:
1. 返回节点名(字符串):
def router(state) -> str:
return "next_node_name"
2. 使用 mapping 映射:
add_conditional_edges("node", router, {
"path_a": "node_a",
"path_b": "node_b",
"end": END,
})
3. 基于多条件路由:
def smart_router(state) -> str:
if state["confidence"] < 0.5:
return "clarify"
if state["retry_count"] >= 3:
return "give_up"
return state["intent"]
4. 基于 LLM 决策路由:
def llm_router(state) -> str:
decision = llm.decide(state["messages"])
return decision # LLM 决定下一步3.7 与其他概念的关联
<- 状态管理:条件路由基于 State 字段判断
-> ReAct Agent:Agent 的 tool_call 检查是条件路由
-> 多 Agent 编排:Supervisor 用条件路由分发任务
4. ReAct Agent 构建
4.1 定义
用 LangGraph 构建 ReAct Agent 是最经典的 Agent 图结构。agent 节点(调用 LLM 决策)和 tool_node(执行工具)之间形成循环,直到模型不再调用工具。核心是 agent 节点 + 条件边检查 tool_calls + tool_node 循环。
对应 Demo: demos/04_LangGraph/04_ReAct_Agent.py
4.2 ReAct Agent 图结构
+------------------------------------------------------------------+
| LangGraph ReAct Agent 图结构 |
+------------------------------------------------------------------+
| |
| START |
| | |
| v |
| +---------------+ |
| | agent | 调用 LLM, 返回 AIMessage |
| | (LLM 决策) | 可能包含 tool_calls |
| +-------+-------+ |
| | |
| v |
| +---------------+ |
| | should_continue| 条件路由 |
| | (检查tool_calls)| |
| +-------+-------+ |
| | |
| +--------+--------+ |
| | | |
| 有 tool_calls 无 tool_calls |
| | | |
| v v |
| +---------------+ END |
| | tool_node | 执行工具, 返回 ToolMessage |
| | (工具执行) | |
| +-------+-------+ |
| | |
| | (回到 agent 继续推理) |
| +-----> agent (循环) |
| |
+------------------------------------------------------------------+4.3 ReAct Agent 实现
class ReActAgentState(TypedDict):
"""ReAct Agent 状态"""
messages: Annotated[list, add_messages] # 消息列表(自动累加)
tool_call_count: Annotated[int, lambda x, y: x + y]
max_tool_calls: int
def agent_node(state: ReActAgentState) -> dict:
"""Agent 节点:调用 LLM 决策"""
messages = state["messages"]
# 真实场景: response = llm.invoke(messages, tools=get_tool_definitions())
response = mock_llm_response(messages, state.get("tool_call_count", 0))
return {"messages": [response]}
def tool_node(state: ReActAgentState) -> dict:
"""工具节点:执行 LLM 请求的工具调用"""
last_message = state["messages"][-1]
tool_messages = []
for tool_call in last_message.tool_calls:
tool_name = tool_call["name"]
tool_args = tool_call["args"]
result = TOOLS[tool_name]["handler"](**tool_args)
tool_messages.append(ToolMessage(
content=result,
tool_call_id=tool_call["id"],
))
return {"messages": tool_messages, "tool_call_count": 1}
def should_continue(state: ReActAgentState) -> str:
"""条件路由:检查是否有 tool_calls"""
last_message = state["messages"][-1]
if hasattr(last_message, "tool_calls") and last_message.tool_calls:
if state.get("tool_call_count", 0) >= state.get("max_tool_calls", 10):
return END # 超过最大调用次数
return "tools"
return END # 无 tool_calls,结束
# 构建 ReAct Agent 图
g = StateGraph(ReActAgentState)
g.add_node("agent", agent_node)
g.add_node("tools", tool_node)
g.add_edge(START, "agent")
g.add_conditional_edges("agent", should_continue, {
"tools": "tools",
END: END,
})
g.add_edge("tools", "agent") # 工具执行后回到 agent
graph = g.compile()4.4 agent <-> tool_node 循环
ReAct 循环执行过程:
Step 1: agent 节点
LLM 输入: [SystemMessage, HumanMessage("查天气")]
LLM 输出: AIMessage(tool_calls=[{name:"search", args:{query:"天气"}}])
Step 2: should_continue
检测到 tool_calls -> 路由到 "tools"
Step 3: tool_node
执行 search(query="天气")
返回: ToolMessage(content="北京晴, 25度")
Step 4: agent 节点 (循环)
LLM 输入: [System, Human, AIMessage(tool_calls), ToolMessage]
LLM 输出: AIMessage(content="北京今天晴天,25度")
Step 5: should_continue
无 tool_calls -> 路由到 END
最终消息列表:
[System, Human, AIMessage(tool_calls), ToolMessage, AIMessage(最终回答)]4.5 防止无限循环
def should_continue(state: ReActAgentState) -> str:
"""带多重保护的停止条件"""
# 1. 检查是否有 tool_calls
last_message = state["messages"][-1]
if not (hasattr(last_message, "tool_calls") and last_message.tool_calls):
return END
# 2. 检查最大工具调用次数
if state.get("tool_call_count", 0) >= state.get("max_tool_calls", 10):
return END
# 3. 检查消息长度(防止上下文溢出)
if len(state["messages"]) > 50:
return END
return "tools"4.6 与其他概念的关联
<- Agent 核心能力 / ReAct:LangGraph ReAct 是 Agent ReAct 的框架实现
-> 持久化与恢复:ReAct Agent 可加 Checkpoint 持久化
-> 人机交互 HITL:ReAct Agent 可加 interrupt_before 实现 HITL
5. 持久化与恢复
5.1 定义
Checkpoint 保存图执行状态,Thread 隔离不同会话,Resume 从断点恢复,Time Travel 回溯历史状态。这是 LangGraph 生产级 Agent 的关键能力--会话中断后可恢复,出错后可回滚,多用户会话互不干扰。
对应 Demo: demos/04_LangGraph/05_持久化与恢复.py
5.2 持久化架构
+------------------------------------------------------------------+
| LangGraph 持久化架构 |
+------------------------------------------------------------------+
| |
| Checkpoint (检查点): |
| -> 每个节点执行后自动保存 State 快照 |
| -> 包含完整状态 + 执行位置 |
| |
| Thread (线程): |
| -> 用 thread_id 隔离不同会话 |
| -> 每个线程有独立的 Checkpoint 序列 |
| |
| Resume (恢复): |
| -> 从最新 Checkpoint 恢复执行 |
| -> 用于会话续接 |
| |
| Time Travel (时间旅行): |
| -> 回溯到任意历史 Checkpoint |
| -> 从历史状态重新执行(修改后重放) |
| |
| Checkpoint 存储: |
| +-- MemorySaver: 内存存储 (开发/测试) |
| +-- SqliteSaver: SQLite 文件 (单机生产) |
| +-- PostgresSaver: PostgreSQL (分布式生产) |
| |
+------------------------------------------------------------------+5.3 Checkpoint 实现
from langgraph.checkpoint.memory import MemorySaver
class ChatState(TypedDict):
"""对话状态"""
messages: Annotated[list, add_messages]
turn_count: Annotated[int, lambda x, y: x + y]
summary: str
def build_persistent_graph():
"""构建带持久化的图"""
g = StateGraph(ChatState)
g.add_node("chat", chat_node)
g.add_edge(START, "chat")
g.add_conditional_edges("chat", should_continue, {
"chat": "chat", # 循环
END: END,
})
# 使用 MemorySaver 作为 Checkpoint 存储
checkpointer = MemorySaver()
return g.compile(checkpointer=checkpointer)5.4 多线程会话隔离
graph = build_persistent_graph()
# 线程 1: 用户 A 的会话
config_a = {"configurable": {"thread_id": "user_a_session_1"}}
graph.invoke({"messages": [HumanMessage("你好")]}, config=config_a)
# -> 保存到 thread_id="user_a_session_1" 的 Checkpoint
# 线程 2: 用户 B 的会话(互不干扰)
config_b = {"configurable": {"thread_id": "user_b_session_1"}}
graph.invoke({"messages": [HumanMessage("hello")]}, config=config_b)
# -> 保存到 thread_id="user_b_session_1" 的 Checkpoint
# 恢复用户 A 的会话
graph.invoke({"messages": [HumanMessage("继续上次的话题")]}, config=config_a)
# -> 从 thread_id="user_a_session_1" 的最新 Checkpoint 恢复多线程隔离示意图:
Checkpoint 存储
+----------------------------------+
| thread_id: "user_a_session_1" |
| checkpoint_1: {messages: [...]} |
| checkpoint_2: {messages: [...]} |
| checkpoint_3: (最新) |
+----------------------------------+
+----------------------------------+
| thread_id: "user_b_session_1" |
| checkpoint_1: {messages: [...]} |
| checkpoint_2: (最新) |
+----------------------------------+
每次调用时通过 thread_id 指定会话
-> 自动从最新 Checkpoint 恢复
-> 新的执行结果追加为新的 Checkpoint5.5 Time Travel(时间旅行)
# 获取所有历史状态
config = {"configurable": {"thread_id": "session_1"}}
history = list(graph.get_state_history(config))
# 查看历史状态
for state in history:
print(f"Step: {state.metadata['step']}")
print(f"Messages: {len(state.values.get('messages', []))}")
# 回溯到特定历史状态并重新执行
target_state = history[3] # 第3个检查点
config_replay = {"configurable": {
"thread_id": "session_1",
"checkpoint_id": target_state.config["configurable"]["checkpoint_id"],
}}
# 从该检查点重新执行
graph.invoke(None, config=config_replay)5.6 Checkpoint 存储对比
5.7 与其他概念的关联
<- ReAct Agent:ReAct Agent 加 Checkpoint 实现持久化
-> 人机交互 HITL:HITL 需要 Checkpoint 暂停和恢复
-> 生产级 Agent:持久化是生产级 Agent 的基础
6. 人机交互 HITL
6.1 定义
LangGraph HITL 用 interrupt_before / interrupt_after 在指定节点暂停,等待人工审批后 Resume。核心能力是工具调用前审批、修改参数、拒绝执行和人工输入反馈,是生产级 Agent 的必备功能。
对应 Demo: demos/04_LangGraph/06_人机交互HITL.py
6.2 HITL 架构
+------------------------------------------------------------------+
| LangGraph HITL 架构 |
+------------------------------------------------------------------+
| |
| START -> agent -> [interrupt_before: tools] -> tools -> agent |
| | |
| v |
| 暂停执行 |
| 保存 Checkpoint |
| | |
| v |
| 人工审批 |
| +---+---+---+ |
| | | | |
| 批准 修改 拒绝 |
| | | | |
| v v v |
| 恢复执行 |
| (从 Checkpoint Resume) |
| |
+------------------------------------------------------------------+6.3 HITL 实现
class HITLState(TypedDict):
"""HITL Agent 状态"""
messages: Annotated[list, add_messages]
tool_call_count: Annotated[int, lambda x, y: x + y]
approval_status: str # pending / approved / rejected
pending_tool: str
def build_hitl_graph():
"""构建带 HITL 的图"""
g = StateGraph(HITLState)
g.add_node("agent", agent_node)
g.add_node("tools", tool_node)
g.add_edge(START, "agent")
g.add_conditional_edges("agent", should_continue, {
"tools": "tools",
END: END,
})
g.add_edge("tools", "agent")
# 关键: 在 tools 节点前中断
checkpointer = MemorySaver()
return g.compile(
checkpointer=checkpointer,
interrupt_before=["tools"], # 工具执行前暂停
)
# 使用 HITL
graph = build_hitl_graph()
config = {"configurable": {"thread_id": "hitl_session"}}
# 第一次调用: 执行到 tools 前暂停
result = graph.invoke(
{"messages": [HumanMessage("发邮件给老板")]},
config=config,
)
# -> Agent 决定调用 send_email,在 tools 前暂停
# 人工查看待执行的工具调用
state = graph.get_state(config)
pending_tool_calls = state.values["messages"][-1].tool_calls
# -> [{"name": "send_email", "args": {"to": "boss@co.com", ...}}]
# 人工决策:
# 选项1: 批准执行
graph.invoke(None, config=config) # 继续执行
# 选项2: 修改参数后执行
graph.update_state(config, {
"messages": [AIMessage(
content="发送邮件",
tool_calls=[{"name": "send_email",
"args": {"to": "manager@co.com", # 修改收件人
"subject": "报告", "body": "数据已查询"},
"id": "modified_call"}],
)]
})
graph.invoke(None, config=config)
# 选项3: 拒绝执行
graph.update_state(config, {"approval_status": "rejected"})
# -> Agent 收到拒绝信号,重新决策6.4 interrupt_before vs interrupt_after
interrupt_before (执行前暂停):
agent -> [暂停] -> tools
-> 用于: 工具执行前审批
-> 人工可以看到将要执行什么,决定是否执行
interrupt_after (执行后暂停):
tools -> [暂停] -> agent
-> 用于: 工具执行后审查
-> 人工可以看到执行结果,决定是否继续
组合使用:
agent -> [before] -> tools -> [after] -> agent
-> 执行前审批 + 执行后审查6.5 HITL 审批流程
完整 HITL 审批流程:
1. Agent 运行到高风险工具前
-> interrupt_before=["tools"] 触发暂停
-> 保存 Checkpoint
2. 系统通知人工
-> 展示: 工具名、参数、风险等级
-> 等待人工决策
3. 人工决策
+-- 批准: graph.invoke(None, config) 继续执行
+-- 修改: graph.update_state() 修改参数后继续
+-- 拒绝: graph.update_state() 标记拒绝,Agent 重新决策
+-- 接管: 人工直接操作,结束 Agent
4. Agent 恢复执行
-> 从 Checkpoint 恢复
-> 根据人工决策继续或调整6.6 与其他概念的关联
<- 持久化与恢复:HITL 依赖 Checkpoint 实现暂停恢复
<- Agent 核心能力 / HITL:LangGraph HITL 是 Agent HITL 的框架实现
-> 生产级 Agent:HITL 是生产级 Agent 的安全组件
7. 子图与并行
7.1 定义
子图(Subgraph)封装复杂逻辑为独立模块,Fan-out/Fan-in 实现并行执行与结果聚合,Map-Reduce 对列表数据批量处理再汇总。这些是构建复杂多步骤工作流的核心模式,让图的层次更清晰、逻辑更模块化。
对应 Demo: demos/04_LangGraph/07_子图与并行.py
7.2 子图架构
+------------------------------------------------------------------+
| 子图与并行模式 |
+------------------------------------------------------------------+
| |
| 1. 子图 (Subgraph): |
| |
| 主图: START -> research_subgraph -> summarize -> END |
| | |
| v |
| +-------+ +-------+ +-------+ |
| |search | -> |analyze| -> | END | |
| +-------+ +-------+ +-------+ |
| (子图内部结构) |
| |
| 2. Fan-out / Fan-in (并行+聚合): |
| |
| +-------+ |
| +--->| 分支A |---+ |
| +-------+ | +-------+ | +-------+ |
| | 分发 |--+ +-->| 聚合 | |
| +-------+ | +-------+ | +-------+ |
| +--->| 分支B |---+ |
| +-------+ |
| |
| 3. Map-Reduce: |
| |
| [item1] -> process -> result1 |
| [item2] -> process -> result2 -> reduce -> final |
| [item3] -> process -> result3 |
| |
+------------------------------------------------------------------+7.3 子图实现
class ResearchState(TypedDict):
"""研究子图状态"""
topic: str
search_results: Annotated[list[str], lambda x, y: x + y]
summary: str
def search_node(state: ResearchState) -> dict:
"""子图节点:搜索"""
return {"search_results": [f"关于 {state['topic']} 的搜索结果"]}
def analyze_node(state: ResearchState) -> dict:
"""子图节点:分析"""
return {"summary": f"基于 {len(state['search_results'])} 条结果的分析"}
def build_research_subgraph():
"""构建研究子图: search -> analyze -> END"""
g = StateGraph(ResearchState)
g.add_node("search", search_node)
g.add_node("analyze", analyze_node)
g.add_edge(START, "search")
g.add_edge("search", "analyze")
g.add_edge("analyze", END)
return g.compile()
# 在主图中使用子图
research_subgraph = build_research_subgraph()
class MainState(TypedDict):
topic: str
research_summary: str
def research_node(state: MainState) -> dict:
"""主图节点:调用子图"""
result = research_subgraph.invoke({"topic": state["topic"]})
return {"research_summary": result["summary"]}7.4 Fan-out / Fan-in 实现
class ParallelState(TypedDict):
"""并行处理状态"""
task: str
results: Annotated[list[str], lambda x, y: x + y] # 累加结果
final_report: str
def branch_a(state: ParallelState) -> dict:
"""并行分支 A"""
return {"results": [f"分支A 处理 {state['task']}"]}
def branch_b(state: ParallelState) -> dict:
"""并行分支 B"""
return {"results": [f"分支B 处理 {state['task']}"]}
def aggregate(state: ParallelState) -> dict:
"""聚合结果"""
return {"final_report": "\n".join(state["results"])}
# 构建并行图
g = StateGraph(ParallelState)
g.add_node("branch_a", branch_a)
g.add_node("branch_b", branch_b)
g.add_node("aggregate", aggregate)
# Fan-out: START 同时指向 branch_a 和 branch_b
g.add_edge(START, "branch_a")
g.add_edge(START, "branch_b")
# Fan-in: branch_a 和 branch_b 都指向 aggregate
g.add_edge("branch_a", "aggregate")
g.add_edge("branch_b", "aggregate")
g.add_edge("aggregate", END)7.5 Map-Reduce 模式
class MapReduceState(TypedDict):
"""Map-Reduce 状态"""
items: list[str]
results: Annotated[list[str], lambda x, y: x + y]
final_result: str
def map_node(state: MapReduceState) -> dict:
"""Map: 对每个 item 并行处理"""
results = [f"处理: {item}" for item in state["items"]]
return {"results": results}
def reduce_node(state: MapReduceState) -> dict:
"""Reduce: 聚合所有结果"""
return {"final_result": f"共处理 {len(state['results'])} 项"}
g = StateGraph(MapReduceState)
g.add_node("map", map_node)
g.add_node("reduce", reduce_node)
g.add_edge(START, "map")
g.add_edge("map", "reduce")
g.add_edge("reduce", END)7.6 模式对比
7.7 与其他概念的关联
<- 图基础:子图是图的嵌套
-> 多 Agent 编排:每个 Agent 可封装为子图
-> 生产级 Agent:复杂工作流需要子图组织
8. 流式输出
8.1 定义
LangGraph stream() 方法支持多种 stream_mode,实时获取执行过程中的中间结果。values 模式输出完整状态快照,updates 模式输出增量更新,debug 模式输出执行轨迹。前端用 SSE(Server-Sent Events)推送给用户。
对应 Demo: demos/04_LangGraph/08_流式输出.py
8.2 流式输出模式
+------------------------------------------------------------------+
| LangGraph 流式输出模式 |
+------------------------------------------------------------------+
| |
| stream_mode="values": |
| 每个节点执行后输出完整 State 快照 |
| [State_0] -> [State_1] -> [State_2] -> [State_final] |
| -> 适合: 需要查看每步完整状态 |
| |
| stream_mode="updates": |
| 每个节点只输出增量更新(返回的 dict) |
| [{node:A, update:{...}}] -> [{node:B, update:{...}}] |
| -> 适合: 只关心变化部分 |
| |
| stream_mode="debug": |
| 输出详细执行轨迹(节点名、输入、输出、时间) |
| -> 适合: 调试和监控 |
| |
+------------------------------------------------------------------+8.3 流式输出实现
class StreamState(TypedDict):
"""流式输出状态"""
messages: Annotated[list, add_messages]
step: str
data: Annotated[list[str], lambda x, y: x + y]
result: str
def step_collect(state: StreamState) -> dict:
"""步骤1:收集数据"""
return {"step": "collected", "data": ["数据1", "数据2", "数据3"]}
def step_process(state: StreamState) -> dict:
"""步骤2:处理数据"""
return {"step": "processed",
"data": [f"已处理:{d}" for d in state["data"]]}
def step_analyze(state: StreamState) -> dict:
"""步骤3:分析结果"""
return {"step": "done", "result": f"分析完成: {len(state['data'])} 项"}
g = StateGraph(StreamState)
g.add_node("collect", step_collect)
g.add_node("process", step_process)
g.add_node("analyze", step_analyze)
g.add_edge(START, "collect")
g.add_edge("collect", "process")
g.add_edge("process", "analyze")
g.add_edge("analyze", END)
graph = g.compile()
# values 模式: 完整状态快照
for event in graph.stream({"messages": [HumanMessage("开始")]},
stream_mode="values"):
print(f"Step: {event.get('step', 'init')}")
print(f"Data: {event.get('data', [])}")
# updates 模式: 增量更新
for event in graph.stream({"messages": [HumanMessage("开始")]},
stream_mode="updates"):
for node_name, update in event.items():
print(f"[{node_name}] 更新: {update}")8.4 输出示例对比
values 模式输出:
{"messages": [HumanMessage("开始")], "step": "", "data": [], "result": ""}
{"messages": [HumanMessage("开始")], "step": "collected", "data": ["数据1","数据2","数据3"]}
{"messages": [HumanMessage("开始")], "step": "processed", "data": ["已处理:数据1",...]}
{"messages": [HumanMessage("开始")], "step": "done", "result": "分析完成: 3项"}
updates 模式输出:
{"collect": {"step": "collected", "data": ["数据1","数据2","数据3"]}}
{"process": {"step": "processed", "data": ["已处理:数据1",...]}}
{"analyze": {"step": "done", "result": "分析完成: 3项"}}
区别:
values: 每次输出完整 State(包含历史)
updates: 每次只输出本节点的更新(更轻量)8.5 SSE 前端推送
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
@app.get("/chat/stream")
def chat_stream(message: str):
"""SSE 端点:流式推送 Agent 执行过程"""
def event_stream():
for event in graph.stream(
{"messages": [HumanMessage(message)]},
stream_mode="updates",
):
for node, update in event.items():
yield f"data: {node}: {update}\n\n"
yield "data: [DONE]\n\n"
return StreamingResponse(event_stream(), media_type="text/event-stream")8.6 流式输出的应用价值
流式输出的价值:
1. 用户体验:
-> 实时看到执行进度,而非等待全部完成
-> "正在搜索..." -> "找到3条结果" -> "正在分析..." -> "完成"
2. 调试监控:
-> 每个节点的输入输出都可见
-> 快速定位问题节点
3. 成本控制:
-> 在中间步骤检测到问题时可以提前终止
-> 不需要等全部执行完才发现错误
4. 前端集成:
-> SSE/WebSocket 推送实时更新
-> 支持取消操作8.7 与其他概念的关联
<- LLM API 调用 / 流式输出:图的流式输出与 LLM 流式输出互补
-> 生产级 Agent:生产环境必须支持流式输出
-> 多 Agent 编排:多 Agent 执行过程的可视化需要流式输出
9. 多 Agent 编排
9.1 定义
多 Agent 编排通过 Supervisor(主管)模式协调多个专家 Agent。Supervisor 负责任务分发和结果聚合,专家 Agent 各司其职。用 LangGraph 的条件路由实现 Supervisor -> 专家 Agent -> Supervisor 的循环,是复杂任务的标准解决方案。
对应 Demo: demos/04_LangGraph/09_多Agent编排.py
9.2 Supervisor 模式架构
+------------------------------------------------------------------+
| 多 Agent 编排 (Supervisor 模式) |
+------------------------------------------------------------------+
| |
| START |
| | |
| v |
| +---------------+ |
| | Supervisor | 分析任务,决定由哪个 Agent 处理 |
| | (主管) | -> next_agent |
| +-------+-------+ |
| | |
| +--------+--------+--------+ |
| | | | | |
| v v v v |
| +------+ +------+ +------+ +------+ |
| |研究 | |编码 | |审查 | |FINISH| |
| |Agent | |Agent | |Agent | | | |
| +--+---+ +--+---+ +--+---+ +------+ |
| | | | |
| +--------+--------+ |
| | |
| v |
| +---------------+ |
| | Supervisor | <- 查看结果,决定下一步 |
| | (回到主管) | |
| +-------+-------+ |
| | |
| v |
| +---------------+ |
| | FINISH | -> 汇总所有结果 |
| +---------------+ |
| | |
| v |
| END |
| |
+------------------------------------------------------------------+9.3 多 Agent 状态
class MultiAgentState(TypedDict):
"""多 Agent 共享状态"""
user_request: str
messages: Annotated[list[str], lambda x, y: x + y]
next_agent: str # Supervisor 决定的下一个 Agent
research_result: str # 研究 Agent 结果
code_result: str # 代码 Agent 结果
review_result: str # 审查 Agent 结果
final_answer: str # 最终答案
agent_history: Annotated[list[str], lambda x, y: x + y]9.4 Supervisor 节点
def supervisor(state: MultiAgentState) -> dict:
"""
Supervisor: 分析任务,决定由哪个专家 Agent 处理
真实场景: 调用 LLM 分析任务并选择 Agent
"""
history = state.get("agent_history", [])
if len(history) == 0:
next_agent = "researcher"
elif "researcher" in history and "coder" not in history:
next_agent = "coder"
elif "coder" in history and "reviewer" not in history:
next_agent = "reviewer"
else:
next_agent = "FINISH"
return {"next_agent": next_agent,
"messages": [f"[Supervisor] 分发 -> {next_agent}"]}
def route_from_supervisor(state: MultiAgentState) -> str:
"""根据 Supervisor 决策路由"""
next_agent = state["next_agent"]
if next_agent == "FINISH":
return "finalize"
return next_agent9.5 专家 Agent 节点
def researcher(state: MultiAgentState) -> dict:
"""研究 Agent: 收集信息"""
return {
"research_result": f"研究完成: {state['user_request']}",
"agent_history": ["researcher"],
"messages": ["[Researcher] 研究完成"],
}
def coder(state: MultiAgentState) -> dict:
"""编码 Agent: 编写代码"""
return {
"code_result": f"代码编写完成 (基于: {state.get('research_result', '')})",
"agent_history": ["coder"],
"messages": ["[Coder] 编码完成"],
}
def reviewer(state: MultiAgentState) -> dict:
"""审查 Agent: 审查结果"""
return {
"review_result": "审查通过",
"agent_history": ["reviewer"],
"messages": ["[Reviewer] 审查完成"],
}
def finalize(state: MultiAgentState) -> dict:
"""汇总结果"""
return {
"final_answer": (f"研究: {state.get('research_result','')}\n"
f"代码: {state.get('code_result','')}\n"
f"审查: {state.get('review_result','')}"),
}9.6 构建多 Agent 图
g = StateGraph(MultiAgentState)
g.add_node("supervisor", supervisor)
g.add_node("researcher", researcher)
g.add_node("coder", coder)
g.add_node("reviewer", reviewer)
g.add_node("finalize", finalize)
g.add_edge(START, "supervisor")
g.add_conditional_edges("supervisor", route_from_supervisor, {
"researcher": "researcher",
"coder": "coder",
"reviewer": "reviewer",
"finalize": "finalize",
})
# 每个 Agent 完成后回到 Supervisor
g.add_edge("researcher", "supervisor")
g.add_edge("coder", "supervisor")
g.add_edge("reviewer", "supervisor")
g.add_edge("finalize", END)
graph = g.compile()9.7 多 Agent 编排模式对比
9.8 Supervisor 决策策略
Supervisor 决策方式:
1. 规则决策 (本 Demo):
根据 agent_history 顺序执行
-> 简单可控,适合固定流程
2. LLM 决策:
调用 LLM 分析当前状态,选择最合适的 Agent
-> 灵活智能,适合动态任务
3. 混合决策:
规则 + LLM
-> 规则处理常见情况,LLM 处理特殊情况
LLM 决策示例:
"当前状态: 研究已完成, 代码已编写
下一步应该做什么?
A. 再次研究 B. 修改代码 C. 审查 D. 完成"
-> LLM 返回 "C"9.9 与其他概念的关联
<- 条件路由与分支:Supervisor 用条件路由分发任务
<- 子图与并行:每个 Agent 可封装为子图
<- Agent 核心能力 / 完整 Agent:多 Agent 是单 Agent 的扩展
-> 生产级 Agent:生产级系统通常需要多 Agent 协作
10. 生产级 Agent
10.1 定义
生产级 LangGraph Agent 集成 ReAct + 工具调用 + Checkpoint 持久化 + HITL 审批 + 流式输出 + 预算控制。这是前面所有 LangGraph 能力的集大成者,展示如何构建一个可部署、可恢复、可监控的生产级 Agent。
对应 Demo: demos/04_LangGraph/10_生产级Agent.py
10.2 生产级 Agent 架构
+------------------------------------------------------------------+
| 生产级 LangGraph Agent 架构 |
+------------------------------------------------------------------+
| |
| +-------------------+ +-------------------+ |
| | State 定义 | | 工具注册表 | |
| | messages | | search (low) | |
| | tool_call_count | | database (medium) | |
| | max_tool_calls | | send_email (high) | |
| | iteration | +--------+----------+ |
| | user_id | | |
| | session_id | v |
| | status | +-------------------+ |
| +---------+---------+ | 预算控制器 | |
| | | max_iterations | |
| v | max_tool_calls | |
| +-------------------+ +--------+----------+ |
| | ReAct 图结构 | | |
| | | v |
| | agent <-- tools | +-------------------+ |
| | | | | | HITL 审批 | |
| | v v | | interrupt_before | |
| | should_continue | | (高风险工具) | |
| +-------------------+ +-------------------+ |
| | |
| v |
| +-------------------+ |
| | Checkpoint 持久化 | |
| | MemorySaver | |
| | (thread_id 隔离) | |
| +-------------------+ |
| | |
| v |
| +-------------------+ |
| | 流式输出 | |
| | stream_mode= | |
| | "updates" | |
| +-------------------+ |
| |
+------------------------------------------------------------------+10.3 生产级 State 定义
class ProductionAgentState(TypedDict):
"""生产级 Agent 状态"""
messages: Annotated[list, add_messages]
tool_call_count: Annotated[int, lambda x, y: x + y]
max_tool_calls: int
max_iterations: int
iteration: Annotated[int, lambda x, y: x + y]
user_id: str
session_id: str
status: str # running / completed / error / hitl_pending10.4 生产级 Agent 构建
def build_production_agent():
"""构建生产级 Agent"""
g = StateGraph(ProductionAgentState)
g.add_node("agent", agent_node)
g.add_node("tools", tool_node)
g.add_edge(START, "agent")
g.add_conditional_edges("agent", should_continue, {
"tools": "tools",
END: END,
})
g.add_edge("tools", "agent")
# 生产级配置
checkpointer = MemorySaver()
return g.compile(
checkpointer=checkpointer,
interrupt_before=["tools"], # HITL: 工具执行前审批
)10.5 生产级 Agent 使用
agent = build_production_agent()
# 1. 创建会话
config = {"configurable": {"thread_id": f"user_{user_id}_session_{session_id}"}}
# 2. 流式执行 + HITL
for event in agent.stream(
{
"messages": [HumanMessage(user_input)],
"max_tool_calls": 10,
"max_iterations": 20,
"user_id": user_id,
"session_id": session_id,
"status": "running",
},
config=config,
stream_mode="updates",
):
for node, update in event.items():
if node == "agent":
# 检查是否需要 HITL
if "__interrupt__" in update:
# 工具执行前暂停,等待审批
state = agent.get_state(config)
pending_tool = state.values["messages"][-1].tool_calls[0]
# 通知前端: 需要人工审批
notify_human(pending_tool)
# 等待人工决策...
# 批准: agent.invoke(None, config=config)
# 拒绝: agent.update_state(config, {"status": "rejected"})
else:
# 流式推送进度
yield f"data: {node}: {update}\n\n"
# 3. 获取最终状态
final_state = agent.get_state(config)10.6 生产级关键能力对照
生产级能力对照:
+---------------------+-------------------+-------------------+
| 能力 | 实现方式 | 价值 |
+---------------------+-------------------+-------------------+
| ReAct 推理 | agent<->tools 循环 | 自主决策 |
| 工具调用 | ToolRegistry | 外部交互 |
| Checkpoint 持久化 | MemorySaver | 会话恢复 |
| 多线程隔离 | thread_id | 多用户支持 |
| HITL 审批 | interrupt_before | 安全保障 |
| 流式输出 | stream_mode | 实时反馈 |
| 预算控制 | max_tool_calls | 成本控制 |
| 错误恢复 | Checkpoint + 重试 | 可靠性 |
| 状态追踪 | get_state() | 可观测性 |
| Time Travel | get_state_history | 调试回溯 |
+---------------------+-------------------+-------------------+10.7 从 LangGraph Agent 到实际部署
部署架构:
用户 -> API Gateway -> FastAPI Server -> LangGraph Agent
|
+-- Checkpoint Storage (PostgreSQL)
+-- Tool Services (搜索/数据库/邮件)
+-- LLM API (OpenAI/Claude)
+-- Monitoring (LangSmith/Datadog)
+-- Alerting (Prometheus/Grafana)
关键部署考量:
1. 并发控制: 限制同时执行的 Agent 数量
2. 超时处理: 设置全局超时和单节点超时
3. 错误恢复: Checkpoint + 自动重试
4. 监控告警: 追踪执行轨迹、Token 消耗、错误率
5. 安全审计: 记录所有工具调用和 HITL 决策
6. 成本控制: 预算限制 + 实时成本追踪
7. 版本管理: 图结构版本化 + A/B 测试10.8 与其他概念的关联
<- 所有前置 LangGraph 能力:生产级 Agent 集成全部能力
<- Agent 核心能力 / 完整 Agent:LangGraph 是完整 Agent 的框架实现
-> LLMOps 生产化:生产级 Agent 需要完整的 LLMOps 体系支撑
概念关系总览
+--------------------------------------------------+
| LangGraph 技术体系 |
+--------------------------------------------------+
基础层 控制层 高级层 生产层
| | | |
v v v v
+-----------+ +---------------+ +---------------+ +-----------+
|图基础 | |条件路由与分支 | |子图与并行 | |持久化与 |
|StateGraph |-->|add_conditional| |Subgraph | |恢复 |
|Node/Edge | |_edges | |Fan-out/Fan-in | |Checkpoint |
+-----+-----+ +-------+-------+ +-------+-------+ |Thread |
| | | +-----+-----+
v v v |
+-----------+ +---------------+ +---------------+ v
|状态管理 | |ReAct Agent | |流式输出 | +-----------+
|覆盖/累加 |-->|agent<->tools | |stream_mode | |人机交互 |
|Reducer | |条件循环 | |values/updates | |HITL |
|add_messages| +-------+-------+ +-------+-------+ |interrupt |
+-----+-----+ | | +-----+-----+
| v v |
| +---------------+ +---------------+ v
| |多 Agent 编排 | |生产级 Agent | +-----------+
| |Supervisor 模式 | |集成全部能力 |<--|预算控制 |
| |专家 Agent 协作 | |可部署可恢复 | |(Agent核心 |
| +---------------+ +---------------+ | 能力) |
| +-----------+
+-----------> (所有层都基于 State 和 Graph) 概念间的依赖关系
图基础 -> 状态管理 -> 条件路由 -> ReAct Agent
|
持久化与恢复 + 人机交互 HITL
|
子图与并行 + 流式输出
|
多 Agent 编排
|
生产级 Agent
评论区