AB
AiBoss
チュートリアル

LangGraph 实战教程:用状态图构建可控的 LLM 智能体工作流

チュートリアル

LangGraph 实战教程:用状态图构建可控的 LLM 智能体工作流

LangGraph 把智能体的执行流程建模成有向图,用 State、Node、Edge、Reducer 四个概念把分支、循环、并行和人工审批都变成可读的代码。本教程从安装配置讲起,逐步覆盖条件路由、并行聚合、动态扇出、自修正循环与 ReAct 工具调用,并给出可直接运行的完整示例。

用脚本手写一个「思考—行动—观察」的循环,是理解智能体内部机制的好办法。但一旦要把它放进真实业务里,问题就来了:多用户会话怎么隔离和持久化?执行到写数据库、发邮件这类危险动作之前,怎么暂停下来等人工确认?失败的任务怎么按指数退避自动重试?能不能回滚到之前某一步改完状态再重跑?多个节点怎么并行、怎么动态展开?这些如果全靠自己造轮子,维护成本会高得离谱。LangGraph 解决的正是这件事:它把智能体的执行流程建模成一张有向图,让分支、循环、并行、中断、回滚都成为代码里显式可见的结构,而不是藏在提示词里的黑盒行为。这篇教程面向已经了解智能体基本原理、准备把工作流做扎实的开发者,从零开始把主要能力过一遍。

准备工作

安装依赖

在终端里执行下面的命令,安装图引擎、LangChain 核心以及模型接入层:

pip install -U langgraph langchain langchain-openai pydantic

配置模型密钥

需要调用大模型的代码,先把 API Key 写进环境变量。Windows PowerShell 与 macOS / Linux 的写法不同:

# Windows (PowerShell)
$env:OPENAI_API_KEY = "your-api-key-here"

# macOS / Linux (Bash)
export OPENAI_API_KEY="your-api-key-here"

在 Python 代码里初始化模型时,用 init_chat_model 可以统一不同厂商的接口,OpenAI、Anthropic 乃至本地 Ollama 模型都能用同一套写法切换:

from langchain.chat_models import init_chat_model

llm = init_chat_model("gpt-4o-mini", model_provider="openai", temperature=0)
# 换用其他厂商时改 model_provider 与模型名即可
# llm = init_chat_model("claude-3-5-sonnet-latest", model_provider="anthropic")

具体可用的模型名、计费方式和配额,以各家官网当前公布的信息为准。

有一点值得先说明:本教程中关于图结构、状态与 Reducer、人工审批流程、时间旅行、重试策略、子图这几部分内容,不依赖任何模型 API Key 也能完整跑通。建议先把它们当作普通 Python 脚本在本地执行一遍,确认每一步的输出符合预期,再接入真实模型。

操作步骤

第一步:理解四个核心概念

LangGraph 的构成要素只有三个,外加一个辅助概念:

  • State:节点之间共享的数据结构,通常用 TypedDict 定义。
  • Node:一个普通 Python 函数,接收 State,返回「想要更新的那部分字段」组成的字典。
  • Edge:节点之间的跳转规则,分为普通的 add_edge 和带条件的 add_conditional_edges
  • Reducer:同一个字段被多次写入时的合并规则——是覆盖,还是追加。

第二步:跑通最小图结构

先不调用模型,只验证图本身能否按预期流转:

from typing import TypedDict
from langgraph.graph import StateGraph, START, END

class SimpleState(TypedDict):
    text: str

def step_a(state: SimpleState) -> dict:
    return {"text": state["text"] + " -> Node A"}

def step_b(state: SimpleState) -> dict:
    return {"text": state["text"] + " -> Node B"}

builder = StateGraph(SimpleState)
builder.add_node("node_a", step_a)
builder.add_node("node_b", step_b)
builder.add_edge(START, "node_a")
builder.add_edge("node_a", "node_b")
builder.add_edge("node_b", END)

graph = builder.compile()
result = graph.invoke({"text": "Start"})
print("执行结果:", result)
# 输出: 执行结果: {'text': 'Start -> Node A -> Node B'}

第三步:分清覆盖与追加

Reducer 是最容易踩坑的地方。State 字段只写类型时,后执行的节点会直接覆盖先前的值;如果希望像日志、消息历史那样不断累积,就要用 Annotated[类型, operator.add] 显式声明合并方式。下面的例子把两种行为放在一起对比:

from typing import Annotated, TypedDict
import operator
from langgraph.graph import StateGraph, START, END

class ReducerState(TypedDict):
    current_step: str
    logs: Annotated[list[str], operator.add]

def worker_1(state: ReducerState) -> dict:
    return {"current_step": "Step 1 完成", "logs": ["worker_1 已处理"]}

def worker_2(state: ReducerState) -> dict:
    return {"current_step": "Step 2 完成", "logs": ["worker_2 已处理"]}

builder = StateGraph(ReducerState)
builder.add_node("w1", worker_1)
builder.add_node("w2", worker_2)
builder.add_edge(START, "w1")
builder.add_edge("w1", "w2")
builder.add_edge("w2", END)

app = builder.compile()
result = app.invoke({"current_step": "初始状态", "logs": ["开始"]})
print("current_step(覆盖):", result["current_step"])
print("logs(追加):", result["logs"])
# current_step(覆盖): Step 2 完成
# logs(追加): ['开始', 'worker_1 已处理', 'worker_2 已处理']

第四步:条件路由

根据输入内容把请求分发到不同处理分支。分叉函数返回一个字符串,图按这个字符串决定下一个节点:

from typing import Literal, TypedDict
from langgraph.graph import StateGraph, START, END

class RouteState(TypedDict):
    query: str
    response: str

def classifier(state: RouteState) -> dict:
    # 实际场景中这里可以调用模型做意图分类
    return {}

def route_by_intent(state: RouteState) -> Literal["tech_support", "billing_support"]:
    if "费用" in state["query"] or "账单" in state["query"]:
        return "billing_support"
    return "tech_support"

def tech_support(state: RouteState) -> dict:
    return {"response": "这里是技术支持,请先检查错误日志。"}

def billing_support(state: RouteState) -> dict:
    return {"response": "这里是账单窗口,可以帮你核对订阅方案。"}

builder = StateGraph(RouteState)
builder.add_node("classifier", classifier)
builder.add_node("tech_support", tech_support)
builder.add_node("billing_support", billing_support)
builder.add_edge(START, "classifier")
builder.add_conditional_edges("classifier", route_by_intent, ["tech_support", "billing_support"])
builder.add_edge("tech_support", END)
builder.add_edge("billing_support", END)

router_app = builder.compile()
print(router_app.invoke({"query": "API 的费用方案是什么", "response": ""}))

第五步:并行执行与结果聚合

从 START 同时引出多条边,这些节点就会并行运行;再用一个列表作为起点连到聚合节点,图会等它们全部完成后再继续:

from typing import Annotated, TypedDict
import operator
from langgraph.graph import StateGraph, START, END

class ParallelState(TypedDict):
    topic: str
    results: Annotated[list[str], operator.add]
    summary: str

def search_web(state: ParallelState) -> dict:
    return {"results": [f"网络检索: {state['topic']} 正在快速增长"]}

def search_internal_db(state: ParallelState) -> dict:
    return {"results": [f"内部数据库: 与 {state['topic']} 相关的项目有 3 个"]}

def aggregate(state: ParallelState) -> dict:
    combined = " / ".join(state["results"])
    return {"summary": f"汇总报告: [{combined}]"}

builder = StateGraph(ParallelState)
builder.add_node("search_web", search_web)
builder.add_node("search_internal_db", search_internal_db)
builder.add_node("aggregate", aggregate)
builder.add_edge(START, "search_web")
builder.add_edge(START, "search_internal_db")
builder.add_edge(["search_web", "search_internal_db"], "aggregate")
builder.add_edge("aggregate", END)

parallel_app = builder.compile()
print(parallel_app.invoke({"topic": "AI 智能体", "results": [], "summary": ""}))

第六步:用 Send API 做动态扇出

当并行分支的数量取决于运行时数据(比如章节数、待处理条目数)时,用 Send 在条件边里动态生成任务:

from typing import Annotated, TypedDict
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send

class ReportState(TypedDict):
    theme: str
    sections: list[str]
    section: str
    contents: Annotated[list[str], operator.add]

def planner(state: ReportState) -> dict:
    return {"sections": ["引言", "架构说明", "总结"]}

def orchestrator(state: ReportState):
    return [Send("worker", {"section": s}) for s in state["sections"]]

def worker(state: ReportState) -> dict:
    return {"contents": [f"## {state['section']}\n内容撰写完成"]}

builder = StateGraph(ReportState)
builder.add_node("planner", planner)
builder.add_node("worker", worker)
builder.add_edge(START, "planner")
builder.add_conditional_edges("planner", orchestrator, ["worker"])
builder.add_edge("worker", END)

app = builder.compile()
res = app.invoke({"theme": "LangGraph", "sections": [], "section": "", "contents": []})
print("生成的章节:\n" + "\n".join(res["contents"]))

第七步:自修正循环

让生成节点和评估节点来回迭代,直到满足质量标准。这里必须给 State 加一个最大尝试次数,否则容易陷入死循环:

from typing import Literal, TypedDict
from langgraph.graph import StateGraph, START, END

class CodeGenState(TypedDict):
    code: str
    test_passed: bool
    attempts: int

def generator(state: CodeGenState) -> dict:
    attempt = state.get("attempts", 0) + 1
    code = "print('Hello World')" if attempt >= 2 else "print('bug')"
    return {"code": code, "attempts": attempt}

def evaluator(state: CodeGenState) -> dict:
    return {"test_passed": "Hello World" in state["code"]}

def check_result(state: CodeGenState) -> Literal["generator", "__end__"]:
    if state["test_passed"] or state["attempts"] >= 3:
        return END
    return "generator"

builder = StateGraph(CodeGenState)
builder.add_node("generator", generator)
builder.add_node("evaluator", evaluator)
builder.add_edge(START, "generator")
builder.add_edge("generator", "evaluator")
builder.add_conditional_edges("evaluator", check_result, ["generator", END])

optimizer_app = builder.compile()
print(optimizer_app.invoke({"code": "", "test_passed": False, "attempts": 0}))
# 输出: {'code': "print('Hello World')", 'test_passed': True, 'attempts': 2}

第八步:接入工具调用的 ReAct 智能体

让模型自行决定何时调用外部工具,并把结果回灌给模型继续推理。内置的 ToolNodetools_condition 已经把这条回路封装好了:

from langchain.chat_models import init_chat_model
from langchain_core.tools import tool
from langgraph.graph import MessagesState, StateGraph, START, END
from langgraph.prebuilt import ToolNode, tools_condition

@tool
def calculate_tax(price: int) -> int:
    """计算指定金额的含税价(按 10% 税率)。"""
    return int(price * 1.10)

tools = [calculate_tax]

llm = init_chat_model("gpt-4o-mini", model_provider="openai", temperature=0)
llm_with_tools = llm.bind_tools(tools)

def agent_node(state: MessagesState) -> dict:
    response = llm_with_tools.invoke(state["messages"])
    return {"messages": [response]}

builder = StateGraph(MessagesState)
builder.add_node("agent", agent_node)
builder.add_node("tools", ToolNode(tools))
builder.add_edge(START, "agent")
builder.add_conditional_edges("agent", tools_condition)
builder.add_edge("tools", "agent")

react_agent = builder.compile()
query = {"messages": [("user", "请计算 5000 元的含税价格")]}
result = react_agent.invoke(query)
print(result["messages"][-1].content)

这里 MessagesState 会自动维护消息列表,tools_condition 负责判断模型是否要求调用工具:要求就转到 tools 节点,否则直接结束。

一个完整示例

把上面的片段串起来,做一个不依赖模型、可直接运行的最小完整流程:先分类意图,再并行收集两路信息,最后汇总输出。

from typing import Annotated, Literal, TypedDict
import operator
from langgraph.graph import StateGraph, START, END

class WorkflowState(TypedDict):
    query: str
    route: str
    findings: Annotated[list[str], operator.add]
    answer: str

def classify(state: WorkflowState) -> dict:
    route = "billing" if "费用" in state["query"] else "tech"
    return {"route": route}

def pick_branch(state: WorkflowState) -> Literal["collect_tech", "collect_billing"]:
    return "collect_billing" if state["route"] == "billing" else "collect_tech"

def collect_tech(state: WorkflowState) -> dict:
    return {"findings": ["技术侧: 建议先查看最近的错误日志"]}

def collect_billing(state: WorkflowState) -> dict:
    return {"findings": ["账单侧: 当前订阅方案可在控制台核对"]}

def summarize(state: WorkflowState) -> dict:
    return {"answer": " | ".join(state["findings"])}

builder = StateGraph(WorkflowState)
builder.add_node("classify", classify)
builder.add_node("collect_tech", collect_tech)
builder.add_node("collect_billing", collect_billing)
builder.add_node("summarize", summarize)
builder.add_edge(START, "classify")
builder.add_conditional_edges("classify", pick_branch, ["collect_tech", "collect_billing"])
builder.add_edge("collect_tech", "summarize")
builder.add_edge("collect_billing", "summarize")
builder.add_edge("summarize", END)

app = builder.compile()
print(app.invoke({"query": "API 的费用方案是什么", "route": "", "findings": [], "answer": ""}))

运行后可以看到 route 被判定为 billingfindings 里出现账单侧的信息,answer 是汇总后的结果。把 classify 换成真实的模型调用、把两个收集节点换成检索或工具调用,就是一个可用的业务骨架。

注意事项

  • Reducer 决定数据是否会被覆盖。只写类型注解的字段,后执行的节点会覆盖先前的值;需要累积的字段必须用 Annotated[类型, operator.add] 声明,否则并行分支的结果会互相冲掉。
  • 循环一定要有退出条件。自修正这类回环结构,务必在 State 里保存尝试次数,并在条件函数里判断上限,否则可能无限循环。
  • 并行分支的汇合点要用列表声明。聚合节点必须写成 add_edge(["节点A", "节点B"], "聚合节点"),图才会等待全部前置节点完成。
  • 动态扇出用 Send,不要硬编码边。分支数量在运行时才确定时,条件边返回 Send 对象列表,而不是固定的节点名。
  • 模型名、可用区域、计费与配额属于易变信息,请以各厂商官网当前公布的内容为准,不要照抄教程里的示例值。
  • 先跑无模型版本。图结构、状态合并、路由、并行、循环这些部分不依赖 API Key,先确认它们的行为正确,再接入模型,排查问题会容易得多。