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

注意:MemorySaverInMemorySaver 的别名,两者完全等价,至今均可使用,并非废弃关系。部分旧教程使用 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 使用! 如果没有配置 checkpointerinterrupt 会静默失效(不会暂停也不会报错),这是新手最常见的踩坑点。

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。

路由

路由

路由

完成

用户请求

Supervisor
主管 Agent

Researcher
研究 Agent

Coder
编码 Agent

Writer
写作 Agent

返回结果

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 内部流程

有 tool_calls

无 tool_calls

START

agent 节点
调用 LLM

tools 节点
执行工具

END

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}")

下一篇:《下篇——生产实战篇》将讲解记忆系统、可观测性、企业级部署、安全最佳实践和完整实战案例。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐