Agent 写着写着就成了意大利面条——LangGraph 焊死有状态图编排(中篇:高级特性篇)
Agent 写着写着就成了意大利面条——LangGraph 焊死有状态图编排(中篇:高级特性篇)
本系列共三篇:上篇(入门概念篇)/ 中篇(高级特性篇)/ 下篇(生产实战篇)。
本文示例基于 LangGraph 稳定版,所有代码均基于最新 API 验证。实际使用时请以 PyPI 和 官方文档 为准。文中 API Key / 密钥均为占位符示例,请勿硬编码到代码或提交到代码仓库。
十一、图的编译与执行
11.1 编译
# 基础编译
app = graph.compile()
# 带 Checkpointer 编译(启用持久化)
from langgraph.checkpoint.memory import InMemorySaver
app = graph.compile(checkpointer=InMemorySaver())
# 带中断点编译(Human-in-the-Loop)
app = graph.compile(
checkpointer=InMemorySaver(),
interrupt_before=["review"] # 在 review 节点前暂停
)
# 带递归限制编译
app = graph.compile(recursion_limit=25) # 最多 25 轮迭代
编译做了什么:验证图结构(有无环、是否可达 END)、构建 Pregel 执行计划、注入 Checkpointer 和中断点。编译后的图是不可变的——不能再加节点或边。
11.2 执行方式
# 1. invoke:同步执行,返回最终结果
result = app.invoke({"messages": [("user", "你好")]})
# 2. stream:流式执行,逐节点返回
for chunk in app.stream({"messages": [("user", "你好")]}):
print(chunk)
# {'node_name': {'messages': [...]}}
# {'next_node': {'answer': '...'}}
# 3. ainvoke / astream:异步执行
result = await app.ainvoke({"messages": [("user", "你好")]})
async for chunk in app.astream({"messages": [("user", "你好")]}):
print(chunk)
# 带 Checkpointer 的执行(需要 thread_id)
config = {"configurable": {"thread_id": "conversation-1"}}
result = app.invoke({"messages": [("user", "你好")]}, config=config)
# 断点恢复:不传 input,从上次中断处继续
result = app.invoke(None, config=config)
11.3 Config(配置参数)
config = {
"configurable": {
"thread_id": "session-001", # 会话 ID(Checkpointer 用)
"user_id": "user-123", # 自定义参数
"model": "gpt-4o", # 动态切换模型
},
"recursion_limit": 50, # 最大递归深度
"tags": ["production", "v2"], # 追踪标签
"metadata": {"env": "prod"}, # 元数据
}
result = app.invoke(inputs, config=config)
thread_id 是最重要的配置——它标识一个会话。同一个 thread_id 的多次调用共享 Checkpointer 中的状态,实现多轮对话和断点恢复。
十二、流式输出(Streaming)
LangGraph 提供三种流式模式,满足不同场景需求。
12.1 节点级流式
每个节点执行完后立即输出更新:
for chunk in app.stream(
{"messages": [("user", "分析 LangGraph 的优势")]},
stream_mode="updates" # 每个节点完成后输出增量
):
for node_name, update in chunk.items():
print(f"[{node_name}] -> {update}")
# [classify] -> {'intent': 'research'}
# [retrieve] -> {'documents': [...]}
# [generate] -> {'answer': '...'}
12.2 消息级流式(Token 流)
LLM 生成时逐 Token 输出,实现打字机效果:
for chunk in app.stream(
{"messages": [("user", "写一首关于秋天的诗")]},
stream_mode="messages" # LLM 逐 token 输出
):
# chunk 是 (message_chunk, metadata) 元组
msg_chunk, metadata = chunk
if msg_chunk.content:
print(msg_chunk.content, end="", flush=True)
# 秋...
# 风起...
# 落叶纷飞...
12.3 自定义流式节点
from langgraph.types import StreamWriter
def streaming_generate(state: State, writer: StreamWriter) -> dict:
"""自定义流式输出节点"""
query = state["query"]
# 分段生成,每段通过 writer 流式输出
sections = ["引言", "分析", "结论"]
for section in sections:
# 模拟 LLM 逐段生成
content = generate_section(section, query)
writer({"section": section, "content": content})
return {"answer": "全部生成完成"}
graph = StateGraph(State)
graph.add_node("generate", streaming_generate)
graph.add_edge(START, "generate")
graph.add_edge("generate", END)
app = graph.compile()
# 使用自定义流式
for chunk in app.stream(
{"query": "分析 AI Agent 趋势"},
stream_mode="custom"
):
print(chunk)
# {'section': '引言', 'content': '...'}
# {'section': '分析', 'content': '...'}
# {'section': '结论', 'content': '...'}
12.4 多模式流式
# 同时获取多种流式数据
for chunk in app.stream(
{"messages": [("user", "搜索 LangGraph")]},
stream_mode=["updates", "messages"], # 同时输出节点更新和 Token 流
subgraphs=True, # 包含子图的事件
):
mode, data = chunk
if mode == "updates":
print(f"[节点更新] {data}")
elif mode == "messages":
msg_chunk, metadata = data
print(msg_chunk.content, end="", flush=True)
十三、持久化与检查点(Checkpointing)
13.1 为什么需要 Checkpointing
没有 Checkpointer 的图是无状态的——每次执行从头开始。有了 Checkpointer:
- 多轮对话:同一 thread_id 的多次调用共享状态
- 断点恢复:服务重启后从上次中断处继续
- Human-in-the-Loop:暂停后人工修改状态再继续
- 时间旅行:回滚到任意历史检查点重新执行
- 调试:查看每一步的 State 快照
13.2 Checkpointer 类型
| Checkpointer | 存储后端 | 适用场景 | 持久性 |
|---|---|---|---|
| InMemorySaver | 内存 | 开发 / 测试 | 进程结束即丢失 |
| SqliteSaver | SQLite 文件 | 单机生产 | 文件持久 |
| PostgresSaver | PostgreSQL | 多机生产 | 数据库持久 |
| RedisSaver | Redis | 高性能场景 | 可配置 TTL |
注意:MemorySaver 是 InMemorySaver 的别名,两者完全等价,至今均可使用,并非废弃关系。部分旧教程使用 MemorySaver,新代码建议使用 InMemorySaver 以保持一致性。
13.3 MemorySaver(内存检查点)
from langgraph.checkpoint.memory import InMemorySaver
checkpointer = InMemorySaver()
app = graph.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "session-1"}}
# 第一轮对话
result1 = app.invoke(
{"messages": [("user", "我叫张三")]},
config=config
)
# 第二轮对话(能记住第一轮的内容)
result2 = app.invoke(
{"messages": [("user", "我叫什么名字?")]},
config=config
)
print(result2["messages"][-1].content) # "你叫张三"
13.4 SQLiteSaver(文件检查点)
from langgraph.checkpoint.sqlite import SqliteSaver
# 同步方式
checkpointer = SqliteSaver.from_conn_string("checkpoints.db")
checkpointer.setup() # 首次使用需初始化表结构
app = graph.compile(checkpointer=checkpointer)
# 异步方式(推荐用于 Web 服务)
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
async with AsyncSqliteSaver.from_conn_string("checkpoints.db") as checkpointer:
await checkpointer.setup()
app = graph.compile(checkpointer=checkpointer)
result = await app.ainvoke(inputs, config=config)
SQLite 文件自动建表,schema 由 LangGraph 管理,不需要手动 migration。适合单机单进程部署。⚠️ 注意:SQLite 不支持多进程/多实例并发写入,多实例 Web 服务请使用 PostgresSaver。
13.5 PostgresSaver(生产环境)
from langgraph.checkpoint.postgres import PostgresSaver
DB_URI = "postgresql://user:***@localhost:5432/langgraph_db"
# 同步方式
checkpointer = PostgresSaver.from_conn_string(DB_URI)
checkpointer.setup() # 自动创建 checkpoints/writes/blobs/migrations 表
app = graph.compile(checkpointer=checkpointer)
# 异步方式
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
async with AsyncPostgresSaver.from_conn_string(DB_URI) as checkpointer:
await checkpointer.setup()
app = graph.compile(checkpointer=checkpointer)
result = await app.ainvoke(inputs, config=config)
PostgresSaver 支持连接池,多个服务实例共享同一个状态库,这是 InMemorySaver 和 SqliteSaver 做不到的。多实例生产环境优先使用 PostgresSaver;单机小体量可以使用 SqliteSaver。
13.6 查看检查点历史
config = {"configurable": {"thread_id": "session-1"}}
# 获取状态历史
for state in app.get_state_history(config):
print(f"Step {state.metadata['step']}: {state.values}")
print(f" Next: {state.next}")
print(f" Config: {state.config}")
# 获取当前状态
state = app.get_state(config)
print(f"当前步骤: {state.next}")
print(f"当前状态: {state.values}")
# 回滚到历史检查点(时间旅行)
history = list(app.get_state_history(config))
if len(history) > 1:
old_state = history[-1] # 最早的检查点
# 用旧 config 重新执行
result = app.invoke(None, config=old_state.config)
生产优化提示:当 State 持续累积大量消息/文档时,checkpoint 存储体积会快速膨胀。生产环境可参考官方最佳实践,考虑对旧会话做 TTL 清理或归档,控制存储成本。
十四、Human-in-the-Loop(人工介入)
⚠️ 关键前提:
interrupt必须搭配 Checkpointer 使用! 如果没有配置checkpointer,interrupt会静默失效(不会暂停也不会报错),这是新手最常见的踩坑点。
14.1 interrupt(中断点)
LangGraph 提供两种中断方式:
方式一:编译时声明中断点
from langgraph.checkpoint.memory import InMemorySaver
app = graph.compile(
checkpointer=InMemorySaver(),
interrupt_before=["execute_tool"], # 在 execute_tool 节点前暂停
interrupt_after=["generate_draft"], # 在 generate_draft 节点后暂停
)
config = {"configurable": {"thread_id": "task-1"}}
# 执行到 execute_tool 前会暂停
result = app.invoke(inputs, config=config)
# 此时图暂停,state 保存了到 execute_tool 之前的所有数据
# 查看当前状态
state = app.get_state(config)
print(f"暂停在: {state.next}") # ['execute_tool']
print(f"当前状态: {state.values}")
# 人工修改状态后继续执行
app.update_state(config, {"approved": True})
# 从中断点继续
result = app.invoke(None, config=config)
方式二:运行时动态中断
from langgraph.types import interrupt, Command
def human_review(state: State) -> dict:
# 在节点内部调用 interrupt,动态暂停
review_request = {
"draft": state["answer"],
"question": "这段内容是否可以发布?"
}
# interrupt 会暂停图执行,返回值给调用方
# 恢复时,interrupt 返回人工输入的数据
human_input = interrupt(review_request)
# 根据人工输入决定下一步
if human_input.get("action") == "approve":
return {"status": "approved", "answer": human_input.get("edited_content", state["answer"])}
else:
return {"status": "rejected", "feedback": human_input.get("feedback", "")}
14.2 完整 Human-in-the-Loop 审批流
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import interrupt
class State(TypedDict):
messages: Annotated[list, add_messages]
draft: str
final: str
status: str
def generate_draft(state: State) -> dict:
draft = f"基于用户请求生成草稿: {state['messages'][-1].content}"
return {"draft": draft}
def human_approve(state: State) -> dict:
# 动态中断,等待人工审批
decision = interrupt({
"type": "approval",
"draft": state["draft"],
"message": "请审核以下内容,回复 approve 或 reject"
})
if decision["action"] == "approve":
return {
"status": "approved",
"final": decision.get("edited_content", state["draft"])
}
else:
return {
"status": "rejected",
"final": "",
"messages": [("user", f"被拒绝: {decision.get('reason', '未知原因')}")]
}
def publish(state: State) -> dict:
if state["status"] == "approved":
return {"messages": [("ai", f"已发布: {state['final']}")]}
return {"messages": [("ai", "内容被拒绝,未发布")]}
# 构建图
graph = StateGraph(State)
graph.add_node("generate", generate_draft)
graph.add_node("approve", human_approve)
graph.add_node("publish", publish)
graph.add_edge(START, "generate")
graph.add_edge("generate", "approve")
graph.add_conditional_edges("approve", lambda s: "publish" if s["status"] == "approved" else END)
graph.add_edge("publish", END)
app = graph.compile(checkpointer=InMemorySaver())
# === 执行 ===
config = {"configurable": {"thread_id": "review-1"}}
# 第一次调用:执行到 approve 节点会暂停
result = app.invoke(
{"messages": [("user", "写一篇关于 AI 的文章")]},
config=config
)
# 查看暂停状态
state = app.get_state(config)
print(f"暂停在: {state.next}")
print(f"草稿内容: {state.values.get('draft')}")
# 人工审批(通过 Command 恢复)
from langgraph.types import Command
result = app.invoke(
Command(resume={"action": "approve", "edited_content": "修改后的内容..."}),
config=config
)
print(result["messages"][-1].content) # 已发布: 修改后的内容...
14.3 获取中断信息
# 获取所有中断点
state = app.get_state(config)
if state.tasks:
for task in state.tasks:
if hasattr(task, 'interrupts'):
for intr in task.interrupts:
print(f"中断类型: {intr.value.get('type')}")
print(f"中断内容: {intr.value}")
print(f"中断节点: {task.name}")
十五、多 Agent 工作流
15.1 Supervisor 模式(主管协调)
一个 Supervisor Agent 统筹多个专业 Agent。Supervisor 决定调用哪个 Agent,Agent 执行完返回结果给 Supervisor。
from typing import TypedDict, Annotated, Literal
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage, AIMessage
llm = ChatOpenAI(model="gpt-4o", temperature=0)
class State(TypedDict):
messages: Annotated[list, add_messages]
next: str
def supervisor(state: State) -> dict:
"""Supervisor 决定下一个 Agent"""
system_prompt = """你是一个团队主管,管理以下专业 Agent:
- researcher: 负责信息搜索和研究
- coder: 负责编写代码
- writer: 负责撰写文档
- FINISH: 任务完成
根据用户需求和当前进度,决定下一步调用哪个 Agent。
只返回 Agent 名称或 FINISH。"""
response = llm.invoke([
HumanMessage(content=system_prompt),
*state["messages"]
])
next_agent = response.content.strip().lower()
return {"messages": [response], "next": next_agent}
def researcher(state: State) -> dict:
response = llm.invoke([
HumanMessage(content="你是研究专家,搜索并分析相关信息。"),
*state["messages"]
])
return {"messages": [response]}
def coder(state: State) -> dict:
response = llm.invoke([
HumanMessage(content="你是编程专家,编写高质量的代码。"),
*state["messages"]
])
return {"messages": [response]}
def writer(state: State) -> dict:
response = llm.invoke([
HumanMessage(content="你是写作专家,撰写清晰的文章。"),
*state["messages"]
])
return {"messages": [response]}
def route_from_supervisor(state: State) -> str:
next_agent = state["next"]
if next_agent == "FINISH":
return END
return next_agent
# 构建图
graph = StateGraph(State)
graph.add_node("supervisor", supervisor)
graph.add_node("researcher", researcher)
graph.add_node("coder", coder)
graph.add_node("writer", writer)
graph.add_edge(START, "supervisor")
graph.add_conditional_edges("supervisor", route_from_supervisor, {
"researcher": "researcher",
"coder": "coder",
"writer": "writer",
END: END,
})
# 所有 Agent 执行完回到 Supervisor
graph.add_edge("researcher", "supervisor")
graph.add_edge("coder", "supervisor")
graph.add_edge("writer", "supervisor")
app = graph.compile()
# 也可以用官方的 langgraph-supervisor 包(⚠️ 实验性库,API 可能变动,生产环境慎用)
# pip install langgraph-supervisor
# from langgraph_supervisor import create_supervisor
# supervisor = create_supervisor(
# agents=[researcher_agent, coder_agent, writer_agent],
# model=llm,
# prompt="你是团队主管..."
# )
# app = supervisor.compile()
15.2 并行 Agent 模式
多个 Agent 同时执行,结果汇总后返回:
class State(TypedDict):
query: str
research_result: str
code_result: str
summary: str
def research_agent(state: State) -> dict:
# 模拟研究
result = f"研究结果: {state['query']}"
return {"research_result": result}
def code_agent(state: State) -> dict:
# 模拟编码
result = f"代码实现: def solve(): ..."
return {"code_result": result}
def summarize(state: State) -> dict:
summary = f"研究: {state['research_result']}\n代码: {state['code_result']}"
return {"summary": summary}
# 并行执行:两个 Agent 从 START 同时出发
graph = StateGraph(State)
graph.add_node("research", research_agent)
graph.add_node("code", code_agent)
graph.add_node("summarize", summarize)
# 并行:两个边都从 START 出发
graph.add_edge(START, "research")
graph.add_edge(START, "code")
# 汇合:两个都完成后执行 summarize
graph.add_edge("research", "summarize")
graph.add_edge("code", "summarize")
graph.add_edge("summarize", END)
app = graph.compile()
# research 和 code 并行执行,都完成后 summarize 执行
# ⚠️ 并行执行语义:LangGraph 中多个节点从同一源节点出发时并行执行,
# 但下游节点必须等待所有上游并行节点全部完成后才会被触发(barrier 同步等待)。
# 如果 research 和 code 都指向 summarize,则 summarize 会等两者都完成才执行。
15.3 反馈循环(Reflection Loop)
Agent 生成内容后自我审查,不满意就重新生成:
class State(TypedDict):
messages: Annotated[list, add_messages]
draft: str
critique: str
iteration: int
quality_score: float
def generate(state: State) -> dict:
draft = llm.invoke([
HumanMessage(content="写一段关于 AI 的科普文章"),
*state["messages"]
])
return {"draft": draft.content, "iteration": state.get("iteration", 0) + 1}
def reflect(state: State) -> dict:
critique = llm.invoke([
HumanMessage(content=f"评价以下文章的质量,给出 0-1 的分数和改进建议:\n{state['draft']}")
])
score = 0.8 # 实际从 critique 中解析
return {"critique": critique.content, "quality_score": score}
def should_continue(state: State) -> str:
if state["quality_score"] >= 0.85:
return "done"
if state["iteration"] >= 3: # 最多迭代 3 次
return "done"
return "rewrite"
graph = StateGraph(State)
graph.add_node("generate", generate)
graph.add_node("reflect", reflect)
graph.add_edge(START, "generate")
graph.add_edge("generate", "reflect")
graph.add_conditional_edges("reflect", should_continue, {
"rewrite": "generate", # 不满意,重新生成(形成循环)
"done": END,
})
app = graph.compile(recursion_limit=10) # 安全阀:最多 10 轮
反馈循环是 LangGraph 图支持环的自然体现——从 reflect 条件边回到 generate,形成自我改进的闭环。
十六、内置 ReAct Agent
16.1 create_react_agent
LangGraph 提供了预置的 ReAct Agent,两行代码创建一个具备工具调用能力的 Agent:
from langgraph.prebuilt import create_react_agent
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
@tool
def search(query: str) -> str:
"""搜索网页信息"""
return f"搜索结果: {query}"
@tool
def calculate(expression: str) -> str:
"""计算数学表达式"""
# ⚠️ 警告:eval() 是高危函数,严禁在生产环境使用。推荐用 numexpr 库替代。
try:
return str(eval(expression))
except Exception as e:
return f"错误: {e}"
llm = ChatOpenAI(model="gpt-4o")
# 两行创建 ReAct Agent
agent = create_react_agent(llm, tools=[search, calculate])
# 执行
result = agent.invoke({
"messages": [("user", "搜索 LangGraph 并计算 2+2")]
})
print(result["messages"][-1].content)
# 带检查点
from langgraph.checkpoint.memory import InMemorySaver
agent = create_react_agent(
llm,
tools=[search, calculate],
checkpointer=InMemorySaver()
)
# 多轮对话
config = {"configurable": {"thread_id": "chat-1"}}
result1 = agent.invoke(
{"messages": [("user", "我叫张三")]},
config=config
)
result2 = agent.invoke(
{"messages": [("user", "我叫什么?")]},
config=config
)
print(result2["messages"][-1].content) # "你叫张三"
# 支持字符串模型标识符(0.2.63+ 新增)
agent = create_react_agent("openai:gpt-4o", tools=[search, calculate])
16.2 ReAct Agent 内部流程
create_react_agent 内部构建了一个 StateGraph:
- State:MessagesState(消息列表)
- agent 节点:调用 LLM,生成 AIMessage 或 tool_calls
- tools 节点:ToolNode,自动执行 tool_calls
- 条件边:LLM 返回有 tool_calls -> 去 tools 节点;没有 -> END
- 循环:tools 执行完回到 agent,形成 ReAct 循环
十七、Tool Calling Agent
17.1 手动构建 Tool Calling Agent
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.prebuilt import ToolNode
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
@tool
def search(query: str) -> str:
"""搜索网页"""
return f"结果: {query}"
@tool
def write_file(filename: str, content: str) -> str:
"""写入文件"""
return f"已写入 {filename}"
tools = [search, write_file]
llm = ChatOpenAI(model="gpt-4o").bind_tools(tools)
class State(TypedDict):
messages: Annotated[list, add_messages]
def agent_node(state: State) -> dict:
response = llm.invoke(state["messages"])
return {"messages": [response]}
def should_use_tools(state: State) -> str:
last_message = state["messages"][-1]
if last_message.tool_calls:
return "tools"
return END
# 构建图
graph = StateGraph(State)
graph.add_node("agent", agent_node)
graph.add_node("tools", ToolNode(tools))
graph.add_edge(START, "agent")
graph.add_conditional_edges("agent", should_use_tools, {
"tools": "tools",
END: END,
})
graph.add_edge("tools", "agent") # 工具执行完回到 agent
app = graph.compile(recursion_limit=25)
result = app.invoke({
"messages": [("user", "搜索 LangGraph 然后把结果写入 result.txt")]
})
十八、错误处理与重试
18.1 节点级重试
from langgraph.pregel import RetryPolicy
# 全局重试策略
retry_policy = RetryPolicy(
max_attempts=3, # 最多重试 3 次
initial_interval=1.0, # 初始等待 1 秒
backoff_factor=2.0, # 指数退避因子
max_interval=10.0, # 最大等待 10 秒
)
app = graph.compile(retry_policy=retry_policy)
# 节点级重试(在 add_node 时指定)
graph.add_node("api_call", api_node, retry=RetryPolicy(max_attempts=5))
18.2 全局错误处理
def fallback_node(state: State) -> dict:
"""全局错误兜底节点"""
return {
"answer": "抱歉,处理过程中出现错误,请重试。",
"error": True,
"messages": [("ai", "服务暂时不可用,请稍后再试。")]
}
def error_router(state: State) -> str:
if state.get("error"):
return "fallback"
return "continue"
# 在关键节点后加错误检查
graph.add_conditional_edges("api_call", error_router, {
"fallback": "fallback",
"continue": "next_node",
})
graph.add_node("fallback", fallback_node)
graph.add_edge("fallback", END)
18.3 超时与递归限制
# 编译时设置递归限制
app = graph.compile(recursion_limit=25) # 默认 25
# 运行时覆盖
result = app.invoke(
inputs,
config={"recursion_limit": 50} # 放宽到 50 轮
)
# 超时控制
import asyncio
try:
result = await asyncio.wait_for(
app.ainvoke(inputs, config=config),
timeout=30.0 # 30 秒超时
)
except asyncio.TimeoutError:
print("Agent 执行超时")
# 从检查点恢复状态
state = app.get_state(config)
递归限制是安全阀——防止 Agent 陷入无限循环。默认 25 轮,生产环境建议根据任务复杂度调整。超时控制是另一层保护,防止单个 LLM 调用卡死整个图。
如果你刚接触 LangGraph,建议先阅读《上篇——入门概念篇》打好基础。
附录 A:核心 API 速查
# === 图构建 ===
from langgraph.graph import StateGraph, START, END, MessagesState
graph = StateGraph(State) # 创建图
graph.add_node(name, fn) # 注册节点
graph.add_edge(A, B) # 普通边
graph.add_conditional_edges(src, router, mapping) # 条件边
graph.compile(checkpointer=, store=, interrupt_before=, interrupt_after=, recursion_limit=) # 编译
# === 执行 ===
app.invoke(inputs, config=) # 同步执行
app.stream(inputs, config=, stream_mode=) # 流式执行
app.ainvoke(inputs, config=) # 异步执行
app.astream(inputs, config=, stream_mode=) # 异步流式
# === 状态管理 ===
app.get_state(config) # 获取当前状态
app.get_state_history(config) # 获取状态历史
app.update_state(config, values) # 更新状态
# === 预置组件 ===
from langgraph.prebuilt import create_react_agent, ToolNode, MessagesState
# === 检查点 ===
from langgraph.checkpoint.memory import InMemorySaver # 内存
from langgraph.checkpoint.sqlite import SqliteSaver # SQLite
from langgraph.checkpoint.postgres import PostgresSaver # PostgreSQL
# === Human-in-the-Loop ===
from langgraph.types import interrupt, Command
interrupt(value) # 动态中断
Command(resume=value) # 恢复执行
Command(update=values) # 更新状态后恢复
# === 流式 ===
stream_mode="updates" # 节点级增量
stream_mode="messages" # Token 级流式
stream_mode="custom" # 自定义流式
stream_mode=["updates", "messages"] # 多模式
# === Store(长期记忆)===
from langgraph.store.memory import InMemoryStore
from langgraph.store.postgres import PostgresStore
store.put(namespace, key, value) # 写入
store.get(namespace, key) # 读取
store.search(namespace, key=, limit=) # 键值搜索(InMemoryStore)
store.search(namespace, query=, limit=) # 向量语义搜索(仅 PostgresStore 支持,InMemoryStore 不支持)
附录 B:快速参考卡片
LangGraph 核心五步:
1. 定义 State(TypedDict + Annotated)
2. 定义节点函数(State -> dict)
3. 创建 StateGraph,add_node 注册
4. add_edge / add_conditional_edges 连接
5. compile(checkpointer=) 编译执行
关键 API:
graph.add_node(name, fn)
graph.add_edge(A, B)
graph.add_conditional_edges(src, router_fn, mapping)
graph.compile(checkpointer=, interrupt_before=)
app.invoke(inputs, config={"configurable": {"thread_id": "xxx"}})
app.stream(inputs, stream_mode="messages")
检查点三件套:
InMemorySaver() -> 开发
SqliteSaver(db) -> 单机生产
PostgresSaver(uri) -> 多机生产
Human-in-the-Loop:
interrupt_before=["node"] -> 编译时声明
interrupt(value) -> 运行时动态
Command(resume=value) -> 恢复执行
附录 D:完整 Hello World(带注释)
"""
LangGraph 完整 Hello World
功能:一个带条件分支和工具调用的简单 Agent
依赖:pip install langgraph langchain-openai
"""
import os
from typing import TypedDict, Annotated, Literal
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.prebuilt import ToolNode
from langchain_openai import ChatOpenAI
from langchain_core.tools import tool
from langchain_core.messages import HumanMessage
# === 1. 配置 ===
os.environ["OPENAI_API_KEY"] = "sk-xxx" # 替换为你的 API Key
# === 2. 定义工具 ===
@tool
def search(query: str) -> str:
"""搜索网页信息"""
return f"搜索结果: {query} 相关信息..."
@tool
def calculate(expression: str) -> str:
"""计算数学表达式"""
# ⚠️ 警告:eval() 是高危函数,严禁在生产环境使用。推荐用 numexpr 库替代。
try:
return str(eval(expression))
except Exception as e:
return f"错误: {e}"
tools = [search, calculate]
# === 3. 初始化 LLM ===
llm = ChatOpenAI(model="gpt-4o", temperature=0)
llm_with_tools = llm.bind_tools(tools)
# === 4. 定义 State ===
class AgentState(TypedDict):
messages: Annotated[list, add_messages] # 消息列表,自动追加
# === 5. 定义节点 ===
def agent_node(state: AgentState) -> dict:
"""Agent 节点:调用 LLM 决定下一步"""
response = llm_with_tools.invoke(state["messages"])
return {"messages": [response]}
def should_use_tools(state: AgentState) -> Literal["tools", END]:
"""条件路由:有工具调用去 tools,否则结束"""
last_message = state["messages"][-1]
if hasattr(last_message, "tool_calls") and last_message.tool_calls:
return "tools"
return END
# === 6. 构建图 ===
graph = StateGraph(AgentState)
# 注册节点
graph.add_node("agent", agent_node)
graph.add_node("tools", ToolNode(tools))
# 连接边
graph.add_edge(START, "agent") # 入口 -> agent
graph.add_conditional_edges("agent", should_use_tools) # agent -> tools 或 END
graph.add_edge("tools", "agent") # tools -> agent(循环)
# === 7. 编译(带检查点)===
app = graph.compile(
checkpointer=InMemorySaver(),
recursion_limit=10 # 安全阀
)
# === 8. 执行 ===
config = {"configurable": {"thread_id": "demo-1"}}
result = app.invoke(
{"messages": [HumanMessage(content="搜索 LangGraph 然后计算 123 * 456")]},
config=config
)
# 打印最终回答
print(result["messages"][-1].content)
# 查看执行历史
for state in app.get_state_history(config):
step = state.metadata.get("step", 0)
next_nodes = state.next
print(f"Step {step}: next={next_nodes}")
下一篇:《下篇——生产实战篇》将讲解记忆系统、可观测性、企业级部署、安全最佳实践和完整实战案例。
更多推荐
所有评论(0)