总结梳理 LangGraph 的基础,涵盖 StateGraph、状态定义 (State Schemas)、节点 (Nodes)、边 (Edges)、Command、Send、Runtime 依赖注入、流式传输 (Streaming) 、工具集成和错误处理。
LangGraph 将智能体工作流建模为有向图 (directed graphs)。其核心架构由以下关键要素构成:
图构建完成后必须调用 builder.compile() 生成可运行对象。
以下示例展示智能助手意图路由 (Router Agent Pattern) 工作流。相比于线性链式调用,动态路由与分支跳转(add_conditional_edges)是体现 LangGraph 图编排能力的机制。
核心流转逻辑:
classify 接收用户输入并判定分类(技术支持 vs 账单服务)。route_intent 根据状态中的分类字段,动态决定控制流走向 tech_support 还是 billing_support 节点。from typing import TypedDict
from langgraph.constants import START, END
from langgraph.graph import StateGraph
# 1. 定义图状态 (TypedDict)
class AgentState(TypedDict):
query: str
category: str
response: str
# 2. 定义节点函数 (仅返回增量更新字典)
def classify_node(state: AgentState) -> dict:
category = "technical" if "代码" in state["query"] else "billing"
return {"category": category}
def tech_node(state: AgentState) -> dict:
return {"response": f"【技术支持】已处理关于“{state['query']}”的技术咨询。"}
def billing_node(state: AgentState) -> dict:
return {"response": f"【账单服务】已查询关于“{state['query']}”的账单明细。"}
# 3. 条件路由函数
def route_intent(state: AgentState) -> str:
return state["category"]
# 4. 构建图与配置动态条件边 (Conditional Edges)
builder = StateGraph(AgentState)
builder.add_node("classify", classify_node)
builder.add_node("tech_support", tech_node)
builder.add_node("billing_support", billing_node)
builder.add_edge(START, "classify")
builder.add_conditional_edges(
"classify",
route_intent,
{
"technical": "tech_support",
"billing": "billing_support",
},
)
builder.add_edge("tech_support", END)
builder.add_edge("billing_support", END)
# 5. 编译与执行
graph = builder.compile()
final_state = graph.invoke({"query": "请问这段 Python 代码报错怎么解决?"})
print(final_state["response"])
State 是贯穿整个图生命周期的全局数据容器。
| 容器类型 | 特点 | 适用场景 |
|---|---|---|
TypedDict (推荐) |
无运行时校验开销,语法轻量,是实战项目中最常采用的方式 | 绝大多数智能体编排与图流转场景 |
Pydantic BaseModel |
支持运行时数据强校验、字段默认值及自定义 Validator | 需要严格校验外部 API 输入输出的场景 |
from typing import TypedDict, List
from typing_extensions import NotRequired
# 可设置 total=False 允许实例化时只传入部分初始字段
class AgentState(TypedDict, total=False):
user_query: str
conversation_history: List[str]
search_results: NotRequired[List[str]] # 可选增量字段
from pydantic import BaseModel, Field
from typing import Optional
class AgentStateModel(BaseModel):
user_id: str = Field(description="用户标识")
query: str = Field(min_length=1, description="用户问题")
final_answer: Optional[str] = None
默认情况下,节点返回同名字段的值会直接覆盖现有状态。如果需要列表追加或特殊合并行为,需使用 Reducer:
| 合并逻辑 | 配置方式 | 说明 |
|---|---|---|
| 覆盖值 (默认) | 无 Annotated 声明 | 直接替换原有键值 |
| 列表追加 | Annotated[list, operator.add] |
将节点返回的列表元素追加至旧列表末尾 |
| 消息合并 | Annotated[list, add_messages] |
内置消息 Reducer,支持按 ID 替换或追加 Message |
from typing_extensions import TypedDict, Annotated
import operator
from langgraph.graph import add_messages
class StateWithReducer(TypedDict):
current_step: str # 默认覆盖
logs: Annotated[list, operator.add] # 列表追加
messages: Annotated[list, add_messages] # 消息自动化合并
节点函数接收当前状态,但必须且仅需返回增量更新字典(Partial Update)。切勿在节点内原地修改全量 state 对象再整体返回,否则易引发并发状态冲突或数据重置。
# 错误做法:直接原地修改并返回整个 state
def wrong_node(state: MyState) -> MyState:
state["final_answer"] = "done"
return state
# 正确做法:仅返回变更字段的增量字典
def correct_node(state: MyState) -> dict:
return {"final_answer": "done"}
节点函数可接收不同的参数签名以满足相应的功能需求:
| 签名格式 | 使用场景 |
|---|---|
def node(state: State) |
仅读取和更新图状态的基础节点 |
def node(state: State, runtime: Runtime[Context]) |
需要依赖注入外部服务(数据库连接、底层客户端、仓库实例)的节点 |
在生产级应用中,建议通过 context_schema 定义依赖类型,并在 invoke / astream 时传入服务实例,避免全局变量依赖:
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.runtime import Runtime
# 1. 定义 Context 依赖类型与 State 类型
class AppContext(TypedDict):
db_service: object
class AppState(TypedDict):
query: str
result: str
# 2. 节点中通过 runtime.context 获取注入的服务
def process_node(state: AppState, runtime: Runtime[AppContext]) -> dict:
db = runtime.context["db_service"]
# 使用 db 服务进行查询...
return {"result": "processed"}
# 3. 绑定 state_schema 与 context_schema
builder = StateGraph(state_schema=AppState, context_schema=AppContext)
builder.add_node("process", process_node)
builder.add_edge(START, "process")
builder.add_edge("process", END)
graph = builder.compile()
# 4. 执行时传入 context 实体
res = graph.invoke(
{"query": "test"},
context={"db_service": "DatabaseInstance"}
)
START:图的唯一入口标识。通过 builder.add_edge(START, "first_node") 配置起跑节点。END:图的终止节点标识。走向 END 即代表该轮图执行完成。add_edge)用于确定性的顺序节点跳转或固定分支并发:
# 顺序执行:node_a -> node_b
builder.add_edge("node_a", "node_b")
# 并行并发:node_a 分发至 node_b 与 node_c
builder.add_edge("node_a", "node_b")
builder.add_edge("node_a", "node_c")
add_conditional_edges)根据运行时状态动态路由至不同节点。
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
class RouteState(TypedDict):
query: str
need_web_search: bool
def route_decision(state: RouteState):
if state["need_web_search"]:
# 返回列表触发并行处理
return ["rag_node", "web_node"]
return "rag_node"
builder = StateGraph(RouteState)
builder.add_node("rag_node", lambda s: {"result": "rag"})
builder.add_node("web_node", lambda s: {"result": "web"})
builder.add_conditional_edges(
START,
route_decision,
{
"rag_node": "rag_node",
"web_node": "web_node"
}
)
在实战(如 SQL 校验与纠错)中,常使用 Lambda 结合回流边构成自动纠错循环:
# 当校验节点失败时流转至纠错节点,纠错完成后重新指向校验节点形成回环
builder.add_conditional_edges(
"validate_node",
lambda state: "execute_node" if state.get("error") is None else "correct_node",
{
"correct_node": "correct_node",
"execute_node": "execute_node"
}
)
# 纠错后重新校验
builder.add_edge("correct_node", "validate_node")
builder.add_edge("execute_node", END)
当图中包含回环重试(如校验失败重新生成)时,可以通过在执行时传入 recursion_limit 限制最大超级步数(Superstep),防止因逻辑缺陷引发死循环:
# 设置最大执行步数为 10 步
result = graph.invoke(
{"query": "input"},
config={"recursion_limit": 10}
)
Command APICommand 可以在节点返回值中同时完成状态更新与显式路由跳转:
from langgraph.types import Command
from typing import Literal
def decision_node(state: dict) -> Command[Literal["node_b", "node_c"]]:
if state.get("score", 0) > 80:
return Command(update={"status": "passed"}, goto="node_c")
return Command(update={"status": "failed"}, goto="node_b")
Send API用于 Map-Reduce 模式下的动态分发(Fan-out),针对输入列表为每个元素动态生成并行节点:
from langgraph.types import Send
from typing import Annotated
import operator
class MapReduceState(TypedDict):
items: list[str]
results: Annotated[list, operator.add]
def map_dispatcher(state: MapReduceState):
# 动态为每个 item 生成一个 worker 节点任务
return [Send("worker_node", {"item": item}) for item in state["items"]]
def worker_node(state: dict) -> dict:
return {"results": [f"Processed {state['item']}"]}
invoke / ainvoke)# 同步调用
final_state = graph.invoke({"query": "hello"})
# 异步调用
final_state = await graph.ainvoke({"query": "hello"})
在生产项目中,推荐将编译后的图实例封装在 Workflow 服务类中,对外统一提供 run() 和 stream() 接口:
class KBQueryWorkflow:
def __init__(self):
# 初始化并编译图对象
self.app = create_query_graph()
async def run(self, initial_state: dict):
return await self.app.ainvoke(initial_state)
async def stream(self, initial_state: dict):
async for output in self.app.astream(initial_state):
yield output
astream)LangGraph 支持通过 .astream 获取执行过程中的流式状态事件:
# 1. updates 模式:接收每一步节点产出的状态增量
async for event in graph.astream(initial_state, stream_mode="updates"):
print(f"Node Update: {event}")
# 2. values 模式:接收每一步执行结束后的全量 State 视图
async for state_snapshot in graph.astream(initial_state, stream_mode="values"):
print(f"State Snapshot: {state_snapshot}")
# 3. messages 模式:捕获大模型节点的流式 Token 输出
async for message_chunk, metadata in graph.astream(initial_state, stream_mode="messages"):
if message_chunk.content:
print(message_chunk.content, end="", flush=True)
ToolNode 是 LangGraph 内置的标准工具执行节点(位于 langgraph.prebuilt)。其核心职责为:
state["messages"] 中最新一条 AIMessage 包含的 tool_calls 参数;ToolMessage 并追加写入状态列表中(需在 State 中对 messages 配置 add_messages Reducer)。配合内置条件边路由 tools_condition,可以构建 ReAct 工具调用循环:
from typing import Annotated
from typing_extensions import TypedDict
from langchain_core.tools import tool
from langchain_openai import ChatOpenAI
from langgraph.graph import StateGraph, START, END, add_messages
from langgraph.prebuilt import ToolNode, tools_condition
# 1. 定义工具
@tool
def multiply(a: int, b: int) -> int:
"""计算两个整数的乘积。"""
return a * b
tools = [multiply]
# 2. 定义状态模式
class State(TypedDict):
messages: Annotated[list, add_messages]
# 3. 绑定工具到模型
llm = ChatOpenAI(model="gpt-4o-mini").bind_tools(tools)
def chatbot_node(state: State) -> dict:
return {"messages": [llm.invoke(state["messages"])]}
# 4. 构建图结构
builder = StateGraph(State)
builder.add_node("chatbot", chatbot_node)
builder.add_node("tools", ToolNode(tools)) # 实例化 ToolNode 节点
builder.add_edge(START, "chatbot")
# 结合 tools_condition 动态判断:若模型产生 tool_calls 则转向 "tools" 节点,否则终止转向 END
builder.add_conditional_edges(
"chatbot",
tools_condition,
{
"tools": "tools",
"__end__": END
}
)
builder.add_edge("tools", "chatbot") # 工具执行完成后返回 chatbot 节点
graph = builder.compile()
ToolNode 与 create_agent 的对比| 维度 | LangChain create_agent (传统 API) |
LangGraph StateGraph + ToolNode (图架构) |
|---|---|---|
| 底层架构 | 基于链式 (Chain) / AgentExecutor 循环 | 基于有向图 (StateGraph),天然支持复杂图拓扑与并行分支 |
| 封装粒度 | 高(内部封装完整的 Model 与 Tool 执行流程) | 低(自由控制节点、边、路由条件及状态更新规则) |
| 状态管理 | 依赖 AgentAction / AgentFinish 传输,状态字段难以自定义扩展 | 自定义 State(支持 TypedDict/Pydantic 及 Reducer 增量合并) |
| 控制流拓展 | 受限(难以在工具执行前后插入自定义计算节点或多路分支) | 极高(可任意添加并行节点、人机协同断点、复杂条件分支与子图) |
| 适用场景 | 快速原型验证、单 Agent 简单工具调用场景 | 生产级复杂工作流、多 Agent 协同系统与复杂业务逻辑编排 |
ToolNode 执行工具 -> 返回结果 -> LLM 根据工具输出继续推理并给出最终回答。interrupt_before=["tools"])暂停图的执行,等待人工审批确认后再推进 ToolNode 运行。ToolNode(handle_tool_errors=True) 可以捕获错误并生成包含报错信息的 ToolMessage。LLM 收到报错信息后可进行自我修正重试,提升系统稳定性。ToolNode,主控制图依据意图识别动态分派任务至不同的工具节点。RetryPolicy)针对网络连接中断、速率限制等网络或瞬时错误,可以在注册节点时配置 RetryPolicy:
from langgraph.types import RetryPolicy
builder.add_node(
"search_node",
search_node_func,
retry_policy=RetryPolicy(
max_attempts=3,
initial_interval=1.0,
retry_on=(ConnectionError, TimeoutError)
)
)
ToolNode)在包含工具调用的 Agent 架构中,推荐使用 ToolNode 捕获工具内部产生的运行时错误并将其封装为 ToolMessage,使模型能够感知错误并自我修正:
from langgraph.prebuilt import ToolNode
# handle_tool_errors=True 会自动将工具执行报错转化为 ToolMessage 节点输出
tool_node = ToolNode(tools=[fetch_user_data], handle_tool_errors=True)
builder.add_node("tools", tool_node)