第1章:LangGraph 简介与核心概念
1.1 什么是 LangGraph
| 概念名称 | 说明 | 注意事项 |
|---|
| LangGraph | 基于 LangChain 构建的图编排框架,用于定义和执行状态化、可控流程的 AI 应用。支持循环、条件跳转、并行等复杂控制流。 | 不是一个独立于 LangChain 的项目,而是其扩展库,需配合 LangChain 使用。 |
| 状态化流程(Stateful Workflow) | 每次执行维护一个共享状态对象,节点可读写该状态,实现信息传递与累积。 | 状态设计需清晰,避免字段命名冲突或过度耦合。 |
| 可视化与可调试性 | 提供 ASCII 图打印和集成 LangSmith 追踪能力,便于理解流程执行路径。 | 初学者应善用 .print_ascii() 快速验证图结构。 |
1.2 LangGraph 与 LangChain 的关系
| 概念名称 | 说明 | 注意事项 |
|---|
| 依赖关系 | LangGraph 是 LangChain 生态的一部分,依赖 langchain-core 和 langchain 包。 | 安装时需确保版本兼容,推荐使用最新稳定版。 |
| 功能定位 | LangChain 提供组件(LLM、Prompt、Tools),LangGraph 提供这些组件的编排逻辑。 | 类比:LangChain 是”零件”,LangGraph 是”装配线”。 |
| 兼容性 | 支持所有 LangChain 的 Runnable 接口组件(如 Chains、Agents、Tools)。 | 自定义节点函数需符合 Runnable 签约(输入字典,输出字典)。 |
| Agent 增强 | LangGraph 可构建更可控的 Agent 流程,替代传统 ReAct 循环中的黑盒执行。 | 可精确控制工具调用顺序、中断条件、回退逻辑。 |
1.3 核心概念:图(Graph)、节点(Node)、边(Edge)、状态(State)
| 概念名称 | 说明 | 注意事项 |
|---|
| 图(Graph) | 由节点和边构成的有向图结构,表示整个执行流程。通常使用 StateGraph 创建。 | 是流程的蓝图,需调用 .compile() 才能执行。 |
| 节点(Node) | 代表一个执行单元,通常是一个函数,接收状态并返回更新后的状态片段。 | 函数必须返回字典,键对应状态字段名。 |
| 边(Edge) | 连接两个节点的路径,决定执行顺序。分为固定边和条件边。 | 边不能形成无法退出的无限循环(除非有意设计)。 |
| 状态(State) | 一个持久的数据结构(如 TypedDict),在节点间传递和累积信息。 | 所有节点共享同一状态对象,更新时采用合并策略(非覆盖)。 |
| 入口点(Start Point) | 图开始执行的第一个节点,通过 .set_entry_point() 设置。 | 必须存在,否则图无法启动。 |
| 终止点(End Point) | 图执行结束的目标,可通过特殊边指向 END 或返回 {"__end__": True}。 | 可有多个终止点,如 success、failure 分支。 |
1.4 支持的执行模式:有向图、条件跳转、循环、并行
| 执行模式 | 说明 | 注意事项 |
|---|
| 有向无环图(DAG) | 节点按固定顺序执行,无循环。最简单的流程形式。 | 适合线性处理任务,如数据清洗 → 分析 → 输出。 |
| 条件跳转 | 根据当前状态决定下一跳节点,通过 add_conditional_edges 实现。 | 条件函数必须返回目标节点名或特殊指令(如 END)。 |
| 循环 | 图中存在回边,可重复执行某节点,如 Agent 的”思考-行动”循环。 | 必须设置中断机制(如最大步数、完成标志),防止死循环。 |
| 并行执行 | 多个节点可同时运行(异步模式下),或通过 fan-out/fan-in 模式批量处理。 | 并行需注意状态并发写入冲突,建议使用不可变更新。 |
| 子图嵌套 | 可将一个图作为另一个图的节点,实现模块化设计。 | 子图需独立编译,且状态结构需与父图兼容。 |
第2章:环境搭建与快速入门
2.1 安装 LangGraph 与依赖
| 操作名称 | 操作细节 | 注意事项 |
|---|
| 安装 LangGraph | pip install langgraph | 安装后自动包含 langchain-core 和其他必要依赖。 |
| 升级 LangChain | pip install --upgrade langchain | 推荐使用 langchain>=0.2.0 以确保兼容性。 |
| 验证安装 | python -c "from langgraph.graph import StateGraph; print('OK')" | 若无报错,则安装成功。 |
| 可选:安装可视化工具 | pip install graphviz && conda install -c conda-forge graphviz(如需导出图像) | ASCII 图无需额外依赖,但导出 PNG/SVG 需 graphviz。 |
| 虚拟环境建议 | 使用 venv 或 conda 创建独立环境 | 避免包版本冲突,便于项目隔离。 |
2.2 第一个 LangGraph 程序:Hello World 流程
| 步骤名称 | 操作细节 | 注意事项 |
|---|
| 导入模块 | from langgraph.graph import StateGraph, END | StateGraph 是最常用的图类型。 |
| 定义状态结构 | class State(TypedDict): message: str | 使用 TypedDict 明确字段类型,便于静态检查。 |
| 定义节点函数 | def say_hello(state: State) -> dict: print("Hello, LangGraph!") return {"message": "Hello, World!"} | 函数参数为状态,返回值为要更新的字段字典。 |
| 创建图对象 | workflow = StateGraph(State) | 图基于状态类构建,后续节点将操作该状态。 |
| 添加节点 | workflow.add_node("hello", say_hello) | 第一个参数是节点名(字符串),第二个是函数。 |
| 设置入口点 | workflow.set_entry_point("hello") | 必须设置,否则编译时报错。 |
| 添加边到结束 | workflow.add_edge("hello", END) | 表示执行完 hello 节点后流程结束。 |
| 编译图 | app = workflow.compile() | 返回一个可调用的 RunnableApp 对象。 |
| 执行流程 | app.invoke({"message": ""}) | 输入初始状态,触发执行,输出最终状态。 |
2.3 运行与调试基本图结构
| 方法/操作 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
invoke() | app.invoke(input_state) | 同步执行图,返回最终状态 | result = app.invoke({"message": ""}) print(result) | 适用于单次调用,阻塞主线程。 |
stream() | app.stream(input_state) | 流式输出每一步的执行事件 | for event in app.stream({"message": ""}): | 事件包含节点名、更新状态等,适合调试和前端实时展示。 |
print_ascii() | app.get_graph().print_ascii() | 打印图的 ASCII 结构图 | app.get_graph().print_ascii() | 非常有用!可直观查看节点连接关系和条件跳转。 |
get_graph() | app.get_graph() | 获取图的 Graphviz 对象 | graph = app.get_graph() graph.draw_mermaid_png(output_file="graph.png") | 需安装 graphviz 才能导出图像。 |
| 调试技巧:添加日志 | 在节点函数中使用 print 或 logging | 查看执行顺序和状态变化 | def debug_node(state): print(f"Current state: {state}") return {...} | 生产环境建议使用回调系统替代 print。 |
第3章:状态管理(State Management)
3.1 定义图状态:State 类与 TypedDict
| 概念/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| State(状态类) | class State(TypedDict): key: type | 定义图中共享状态的结构,字段名和类型需提前声明。 | from typing import TypedDict class State(TypedDict): message: str step_count: int finished: bool | 必须继承 TypedDict,字段为可选时可用 total=False。 |
| 字段可选性 | class State(TypedDict, total=False): optional_key: str | 允许某些字段在初始状态中不存在。 | class State(TypedDict, total=False): user_input: str processed: bool | 若字段非 total=False,则所有节点访问时必须存在。 |
| 使用 Pydantic 模型(替代方案) | from langchain_core.pydantic_v1 import BaseModel | 使用 Pydantic 模型作为状态结构,支持验证和默认值。 | class State(BaseModel): message: str = "" count: int = 0 | 需确保使用 pydantic_v1(LangChain 兼容版本),不推荐用于高性能场景。 |
3.2 状态的更新与合并策略
| 概念/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 状态更新 | return {"field": new_value} | 节点函数返回一个字典,表示对状态的增量更新。 | def update_msg(state): return {"message": "updated"} | 返回值不会覆盖整个状态,仅合并指定字段。 |
| 合并策略(Merge) | 自动合并 | LangGraph 自动将返回字典与当前状态按键合并。 | 初始状态: {"msg": "hi", "cnt": 1} 节点返回: {"msg": "hello"} 结果: {"msg": "hello", "cnt": 1} | 若字段值为列表或字典,不会深度合并,而是直接替换。 |
| 列表追加示例 | return {"history": state["history"] + [new_item]} | 手动实现列表累积,避免直接替换。 | def add_to_history(state): return {"history": state["history"] + ["user said hello"]} | 若直接返回新列表,会丢失旧数据。 |
| 字典字段更新 | return {"metadata": {**state["metadata"], "key": "value"}} | 实现字典字段的浅层合并。 | def update_meta(state): return {"metadata": {**state.get("metadata", {}), "source": "web"}} | 建议使用 get 防止 KeyError。 |
3.3 可变 vs 不可变状态设计
| 设计模式 | 说明 | 代码示例 | 注意事项 |
|---|
| 可变状态(Mutable) | 直接修改传入的状态对象。 | def bad_node(state): state["count"] += 1 # 直接修改 return {} | 不推荐:可能导致副作用、难以调试、并发问题。 |
| 不可变设计(Immutable) | 始终返回新值,不修改输入状态。 | def good_node(state): return {"count": state["count"] + 1} | 推荐做法:函数纯净,易于测试和推理。 |
| 深拷贝避免污染 | 当必须修改嵌套结构时,先深拷贝。 | import copy def process_nested(state): new_data = copy.deepcopy(state["data"]) new_data["flag"] = True return {"data": new_data} | — |
| 函数纯净性 | 节点函数应无副作用 | def pure_node(state): result = compute(state["input"]) return {"output": result} | 避免 print、文件写入、全局变量修改等。 |
第4章:节点(Nodes)与边(Edges)
4.1 定义节点函数:函数签名与返回值
| 要素 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 基本签名 | def node_func(state: Dict) -> dict: | 接收当前状态,返回要更新的字段。 | def greet(state): return {"message": "Hello!"} | 参数名不限,但类型建议标注。 |
| 异步节点 | async def async_node(state) -> dict: | 支持异步操作(如 API 调用)。 | async def fetch_data(state): data = await http.get("/api") return {"data": data} | 需使用 ainvoke 或 astream 调用。 |
| 接收配置参数 | def node_with_config(state, config): | 获取运行时配置(如 metadata、callbacks)。 | def log_step(state, config): print(f"Step {config['configurable']['step_id']}") return {} | config 是标准参数,可用于日志、追踪等。 |
| 工具调用集成 | def tool_node(state): return {"result": tool.invoke(state["query"])} | 在节点中调用 LangChain Tool。 | from langchain_core.tools import tool @tool def search(query: str) -> str: return "result" def call_search(state): return {"result": search.invoke(state["query"])} | 工具需提前定义并可调用。 |
| 返回特殊指令 | return {"end": True} 或 raise StopIteration | 主动终止流程(较少用)。 | def check_done(state): if state["done"]: return {"end": True} return {} | 更推荐通过条件边跳转到 END。 |
4.2 添加节点到图:add_node 方法详解
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
add_node | graph.add_node(node_name, node_function) | 将一个函数注册为图中的执行节点。 | workflow = StateGraph(State) workflow.add_node("greet", greet) workflow.add_node("process", process) | node_name 必须唯一,否则覆盖。 |
| 节点命名规范 | 使用小写字母、下划线 | 提高可读性和一致性。 | "validate_input", "call_tool" | 避免空格、特殊字符。 |
| 添加多个节点 | 分别调用 add_node | 逐一注册所有需要的节点。 | workflow.add_node("step1", f1) workflow.add_node("step2", f2) | 顺序不影响执行逻辑,由边决定。 |
| 节点复用 | 同一函数可绑定多个节点名 | 实现相同逻辑的不同上下文调用。 | workflow.add_node("validate_user", validate) workflow.add_node("validate_admin", validate) | 函数内部可通过 config 区分上下文。 |
| 动态添加节点 | 在 compile 前任意时刻添加 | 支持运行时构建图结构。 | if condition: workflow.add_node("special", special_func) | 必须在 .compile() 之前完成所有添加。 |
4.3 边的类型:固定边、条件边、动态边
| 边类型 | 说明 | 适用场景 | 注意事项 |
|---|
| 固定边(Fixed Edge) | 从一个节点无条件跳转到另一个节点。 | 线性流程、必经步骤。 | 使用 add_edge(from, to) 添加。 |
| 条件边(Conditional Edge) | 根据节点返回值或状态决定下一跳。 | 分支判断、if-else 逻辑。 | 需通过 add_conditional_edges 配置条件函数。 |
| 动态边(Dynamic Edge) | 在运行时根据计算结果选择目标节点(实验性)。 | 复杂路由、插件式流程。 | 当前 LangGraph 主要通过条件边模拟此行为。 |
| 到 END 的边 | 显式结束流程。 | 成功、失败、超时等终止路径。 | 目标可为字符串 "END" 或常量 END。 |
| 自循环边 | 节点跳转回自己。 | 实现重复执行(如重试、等待)。 | 必须有退出条件,否则无限循环。 |
4.4 添加边:add_edge 与 add_conditional_edges
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
add_edge | graph.add_edge(start_node, end_node) | 添加一条固定边,执行完 start_node 后跳转到 end_node。 | workflow.add_edge("greet", "process") workflow.add_edge("process", END) | end_node 可为其他节点或 END。 |
add_conditional_edges | graph.add_conditional_edges(source: str, path: Callable, [mapping: Dict]) | 根据 path 函数返回值决定下一跳节点。 | def route_based_on_score(state): if state["score"] > 80: return "approve" else: return "reject" workflow.add_conditional_edges("evaluate", route_based_on_score) | path 函数必须返回有效的节点名或 END。 |
| 条件函数参数 | def condition(state) -> str: | 接收当前状态,返回目标节点名。 | def decide_next(state): return "step_a" if state["flag"] else "step_b" | 不可抛出异常,否则流程中断。 |
| 使用映射表(mapping) | 可选参数,将条件函数返回的键映射到节点名。 | 提高可读性,分离逻辑与路由。 | workflow.add_conditional_edges("classify", classify_func, {"news": "fetch_news", "weather": "get_weather"}) | — |
| 添加多条固定边 | 多次调用 add_edge | 构建线性或分叉结构。 | workflow.add_edge("A", "B") workflow.add_edge("B", "C") workflow.add_edge("C", END) | 边的添加顺序不影响图逻辑。 |
第5章:构建与编译图(Graph Construction & Compilation)
5.1 初始化图对象:StateGraph 与 MessageGraph
| 图类型 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
StateGraph | StateGraph(State) | 基于自定义状态结构(TypedDict 或 Pydantic)构建通用状态化流程图。 | from langgraph.graph import StateGraph class State(TypedDict): message: str step: int workflow = StateGraph(State) | 最常用类型,适用于大多数需要共享状态的场景。 |
MessageGraph | MessageGraph() | 专为消息流设计,默认状态包含 messages 列表,节点操作消息列表。 | from langgraph.graph import MessageGraph workflow = MessageGraph() workflow.add_node("speak", lambda state: {"messages": ["Hello"]}) | 适合对话系统、聊天机器人等以消息为中心的应用。 |
| 状态字段自动处理 | messages 字段自动追加 | MessageGraph 对 messages 字段采用”追加”策略而非替换。 | def respond(state): return {"messages": ["Hi back!"]} # 自动追加到列表 | 若返回非列表值会报错,必须返回消息列表片段。 |
| 初始化参数 | StateGraph(state_schema) | 必须传入状态类(TypedDict 子类)作为 schema。 | class State(TypedDict): ... graph = StateGraph(State) | 不传或类型错误将导致运行时异常。 |
| 图对象可变性 | 可持续添加节点和边 | 图对象在编译前是可变的,支持动态构建。 | if debug_mode: workflow.add_node("debug", debug_fn) | 所有结构修改必须在 .compile() 之前完成。 |
5.2 设置入口点(entry_point)
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
set_entry_point | graph.set_entry_point(node_name) | 指定图开始执行的第一个节点。 | workflow.set_entry_point("start") | 必须设置,否则调用 .compile() 时会抛出错误。 |
| 入口点唯一性 | 每个图仅支持一个入口点 | 流程从单一节点启动。 | workflow.set_entry_point("input_validation") | 不支持多入口(如 fan-in 起始),需通过条件逻辑模拟。 |
| 入口点节点必须存在 | 添加边前需确保节点已注册 | 防止引用未定义节点。 | workflow.add_node("init", init_fn) workflow.set_entry_point("init") | 若节点未添加就设为入口点,编译时报错。 |
| 动态设置入口点 | 可在条件分支后重新设置(不推荐) | 一般不用于运行时切换。 | # 不推荐做法 if env == "test": workflow.set_entry_point("mock_start") | 应通过状态字段控制流程,而非修改图结构。 |
5.3 设置终止点(finish_point)
| 概念/方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 显式终止(END) | 使用 END 常量或字符串 "__end__" | 表示流程正常结束。 | from langgraph.graph import END workflow.add_edge("done", END) | 推荐使用 END 常量以提高可读性。 |
| 多个终止点 | 不同分支可指向 END | 支持成功、失败、超时等多种结束路径。 | workflow.add_edge("success", END) workflow.add_edge("failure", END) | 合法且常见设计。 |
| 返回特殊字段终止 | return {"end": True} | 在节点函数中主动终止(不推荐)。 | def check_end(state): if state["done"]: return {"end": True} return {} | 更推荐通过条件边跳转到 END。 |
| 无终止点风险 | 图中无任何路径可达 END | 导致流程无法正常退出。 | 错误示例:A → B → C → A(循环但无出口) | 必须确保所有路径最终可结束或有中断机制。 |
| 中断与取消 | raise CancelledError 或外部中断 | 支持异步取消执行。 | async with asyncio.timeout(10): await app.ainvoke(...) | 适用于长时间运行的流程。 |
5.4 编译图:compile() 方法与中间件支持
| 方法/参数 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
compile() | graph.compile() | 将图结构编译为可执行的 RunnableApp 对象。 | app = workflow.compile() | 必须调用此方法才能执行流程。 |
checkpointer 参数 | compile(checkpointer=...) | 启用检查点,支持持久化和恢复。 | from langgraph.checkpoint.memory import MemoryCheckpoint app = workflow.compile(checkpointer=MemoryCheckpoint()) | 若启用,每次节点执行后自动保存状态。 |
interrupt_before | compile(interrupt_before=[nodes]) | 在指定节点前暂停执行,等待外部输入。 | app = workflow.compile(interrupt_before=["human_review"]) | 适用于需要人工审批的场景。 |
interrupt_after | compile(interrupt_after=[nodes]) | 在指定节点执行后暂停。 | app = workflow.compile(interrupt_after=["data_fetched"]) | 可用于审查输出或注入额外数据。 |
debug 模式 | compile(debug=True) | 启用调试模式,输出详细执行日志。 | app = workflow.compile(debug=True) | 有助于开发阶段排查问题。 |
| 中间件支持 | compile() 返回对象支持 callbacks | 可传入回调函数监听事件。 | app.invoke(input, config={"callbacks": [...]}) | 与 LangChain 回调系统集成,支持日志、监控等。 |
| 编译后不可变 | 编译后的图结构固定 | 不能再添加节点或边。 | # 错误: app.add_node("new", func) # AttributeError | 所有结构定义必须在 compile 前完成。 |
第6章:条件边与控制流
6.1 条件函数设计:返回边名称
| 要素 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 条件函数签名 | def condition(state) -> str: | 接收当前状态,返回下一跳节点名。 | def route_by_type(state): if state["type"] == "A": return "process_a" else: return "process_b" | 必须返回字符串,表示目标节点名或 END。 |
| 返回值要求 | 字符串,必须匹配已有节点名 | 决定流程跳转路径。 | return "validate" | 若返回无效名称,运行时报错 InvalidNextNode。 |
| 不可抛出异常 | 函数应处理所有情况 | 防止流程意外中断。 | def safe_route(state): return state.get("next", "default") | 建议使用 .get() 和默认值。 |
| 无参数调用 | 条件函数只接收状态 | 不支持额外参数。 | # 不能写成: def cond(state, flag): ... | 如需配置,可通过 config 或闭包传递。 |
| 异步条件函数 | async def async_condition(state) -> str: | 支持异步判断逻辑。 | async def check_db(state): exists = await db.has_record(state["id"]) return "update" if exists else "create" | 需使用 ainvoke 触发异步执行。 |
6.2 多分支跳转与默认边(default route)
| 方法/概念 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 多分支条件函数 | if-elif-else 结构 | 实现多个判断路径。 | def route_task(state): if state["task"] == "write": return "writer" elif state["task"] == "review": return "reviewer" else: return "fallback" | 覆盖所有可能状态值。 |
| 默认边(Default Route) | 在条件函数末尾返回默认节点 | 处理未匹配情况。 | else: return "fallback_handler" | 提高鲁棒性,避免无路由可走。 |
| 使用 mapping 参数设置默认 | add_conditional_edges(..., mapping, default=...) | 显式指定默认目标。 | workflow.add_conditional_edges("classify", classify_fn, mapping={...}, default="unknown") | 若 mapping 中无对应键,则跳转到 default 节点。 |
default 为必需节点 | 默认目标节点必须存在 | 防止路由失败。 | workflow.add_node("unknown", handle_unknown) workflow.add_conditional_edges(..., default="unknown") | 否则运行时报错。 |
| 避免遗漏分支 | 使用 else 或 default | 确保所有状态都有处理路径。 | # 推荐: return routing_map.get(state["key"], "default") | 提升代码健壮性。 |
6.3 使用 Enum 或字符串作为路由键
| 方式 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 使用字符串常量 | 直接返回节点名字符串 | 简单直接。 | return "approval_pending" | 易出错,拼写错误难发现。 |
| 使用 Enum 枚举 | 定义 Enum 类统一管理路由键 | 提高可维护性和类型安全。 | from enum import Enum class NextStep(str, Enum): APPROVE = "approve" REJECT = "reject" def route(state): return NextStep.APPROVE.value | 推荐用于大型项目,避免魔法字符串。 |
| Enum 与 mapping 结合 | 将 Enum 成员映射到节点 | 分离逻辑与配置。 | workflow.add_conditional_edges("judge", get_verdict, mapping={Verdict.PASS: "success", Verdict.FAIL: "retry"}) | 清晰解耦,便于扩展。 |
| 类型提示辅助 | 为条件函数添加返回类型 | 提高可读性。 | def decide(state) -> Literal["A", "B", "END"]: ... | 配合 IDE 提供自动补全和检查。 |
6.4 基于状态字段的条件判断
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 直接访问字段 | state["field"] | 获取状态中特定字段值进行判断。 | def is_done(state): return "end" if state["done"] else "continue" | 确保字段存在,否则 KeyError。 |
| 安全访问字段 | state.get("field", default) | 避免 KeyError,提供默认值。 | def has_data(state): return "process" if state.get("data") else "fetch" | 推荐做法,尤其字段可能不存在时。 |
| 嵌套字段判断 | state["outer"]["inner"] | 访问嵌套结构。 | def check_status(state): return "ok" if state["user"]["active"] else "block" | 需确保路径存在,否则报错。 |
| 安全嵌套访问 | 使用辅助函数或 try-except | 防止深层访问出错。 | def safe_get(d, *keys): for k in keys: if isinstance(d, dict) and k in d: d = d[k] else: return None return d | 复杂场景建议封装安全访问逻辑。 |
| 基于列表长度判断 | len(state["items"]) > 0 | 根据集合大小决定流程。 | def has_items(state): return "process" if state["items"] else "empty" | 常用于处理批量任务或消息队列。 |
第7章:循环与中断控制
7.1 如何实现循环执行
| 方法 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 条件边回跳 | add_edge("node_a", "node_b") & add_conditional_edges(...) | 使用条件函数让流程返回到之前节点,形成循环。 | def loop_condition(state):
return "process" if state["counter"] < 5 else "end"
workflow.add_conditional_edges("process", loop_condition, mapping={"process": "process", "end": END}) | 需确保有终止条件,避免无限循环。 |
| 自定义计数器或状态字段 | 在状态中维护一个计数器或其他状态字段 | 用于跟踪循环次数或满足退出条件的逻辑。 | def increment_counter(state):
state["counter"] += 1
return {"counter": state["counter"]}
workflow.add_node("increment", increment_counter) | 状态字段应在初始化时设置默认值。 |
7.2 防止无限循环:最大步数限制
| 方法/参数 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 设置最大步数 | compile(max_steps=...) | 在编译图时设置最大执行步数,超过该步数将自动结束。 | app = workflow.compile(max_steps=100) | 过低的 max_steps 可能导致正常流程被错误终止。 |
捕获 MaxStepsExceeded 异常 | 使用 try-except 结构处理超步异常 | 当达到最大步数时,可以捕获异常并进行相应处理。 | try:
app.invoke(input_state)
except MaxStepsExceeded:
print("流程因达到最大步数而终止") | 处理程序应设计成安全恢复或记录状态以便后续分析。 |
7.3 主动中断流程:RETURN 与 END 特殊指令
| 指令 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
RETURN | 返回特定键值对(如 {"__return__": value}) | 中断当前流程,并直接返回给调用者指定的结果。 | def early_exit(state):
if state["should_exit"]:
return {"__return__": "Exited early"} | 应确保在需要的地方正确检查和处理 __return__ 字段。 |
END | 使用 END 常量或字符串 "__end__" | 正常结束流程,不同于 RETURN,不返回特定值。 | workflow.add_edge("final_step", END) | END 是更常用的终止方式,适合于流程自然结束。 |
7.4 使用 should_interrupt 控制执行器
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
should_interrupt 回调 | 定义 should_interrupt(state) 函数 | 允许外部逻辑决定是否中断流程。 | def check_interrupt(state):
return state.get("interrupt_signal", False)
app = workflow.compile(should_interrupt=check_interrupt) | 需要在状态中提供某种机制来设置中断信号。 |
| 动态决策中断 | 根据实时数据或外部事件动态决定 | 提供了灵活性,但增加了复杂度。 | async def async_check_interrupt(state):
status = await fetch_status()
return status == "STOP_REQUESTED" | 对于异步操作,需考虑超时和重试策略。 |
第8章:高级图结构
8.1 子图(Subgraphs)与嵌套图
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 创建子图 | SubGraph(name) 或继承自 StateGraph/MessageGraph | 组织相关节点为独立模块,便于管理和复用。 | sub_workflow = StateGraph(State)
sub_workflow.add_node("step1", step1_fn) | 子图需单独编译后才能集成到主图中。 |
| 将子图添加为主图的一部分 | main_graph.add_subgraph(subgraph, entry_point) | 把子图作为主图的一个节点。 | main_workflow.add_subgraph(sub_workflow, "start_sub") | 子图的入口点必须明确定义。 |
| 通过子图减少重复 | 利用子图封装通用逻辑 | 提高代码重用性,简化维护工作。 | validation_subgraph = create_validation_subgraph()
main_workflow.add_subgraph(validation_subgraph, "validate_input")
main_workflow.add_subgraph(validation_subgraph, "validate_output") | 注意子图之间的状态传递和隔离。 |
8.2 并行执行:ParallelNodes 与 Fan-out/Fan-in
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
ParallelNodes | ParallelNodes([nodes]) | 启动一组节点并行执行。 | parallel_tasks = ParallelNodes(["task1", "task2", "task3"]) | 所有节点必须能够独立运行且互不影响。 |
| Fan-out/Fan-in 模式 | 分支执行多个任务然后合并结果 | 支持并发任务处理后的结果聚合。 | def merge_results(states):
return {"merged": states}
workflow.add_parallel_nodes("tasks", ["taskA", "taskB"], merge_results) | 设计良好的合并逻辑至关重要,以确保最终状态一致性。 |
| 并行执行的挑战 | 资源竞争、同步问题等 | 需要特别注意并发环境下的数据一致性和线程安全。 | def safe_task(state):
with lock:
... | 使用锁或其他同步机制保护临界区。 |
8.3 异常处理与错误边(error edges)
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 错误边 | add_error_edge(node_name, error_handler) | 当节点抛出异常时,跳转至指定错误处理器。 | def handle_error(state, exc):
return {"error": str(exc)}
workflow.add_error_edge("risky_operation", "handle_error") | 每个节点可配置一个或多个错误处理器。 |
| 自定义异常处理器 | 实现特定异常处理逻辑 | 提供更加灵活的错误响应策略。 | class CustomException(Exception): pass
def custom_handle(state, exc):
if isinstance(exc, CustomException):
return {"handled_custom": True}
raise exc | 处理器应当清晰区分不同类型的异常。 |
使用 TryNode 包装节点 | TryNode(node, on_error=...) | 为单个节点包裹一层异常捕获机制。 | risky_node = TryNode("risky", on_error="recover") | 相比手动添加错误边,这种方式更为简洁。 |
8.4 动态图构建:运行时修改结构(实验性)
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 动态添加节点/边 | 编译后通过 API 修改图结构 | 支持根据运行时信息调整流程。 | app.add_node("new_step", new_step_fn)
app.add_edge("existing", "new_step") | 当前版本可能有限制,需查阅文档确认兼容性。 |
| 动态特性限制 | 可能影响性能及稳定性 | 动态修改需谨慎,避免引入复杂性。 | if dynamic_condition(state):
app.add_dynamic_path("conditional_step") | 不推荐频繁修改,理想情况下应在编译阶段确定完整结构。 |
| 实验性功能 | 特性可能变更或移除 | 动态图构建仍处于探索阶段,未来可能会发生变化。 | 关注官方更新日志和文档变化 | 依赖此功能的应用应注意版本兼容性。 |
第9章:持久化与检查点(Checkpoints)
9.1 检查点的作用与生命周期
| 概念 | 描述 | 注意事项 |
|---|
| 检查点作用 | 记录流程执行到某一时刻的状态,以便在重启或错误后恢复。 | 定期创建检查点以减少数据丢失风险。 |
| 生命周期 | 从创建检查点开始,直到流程正常结束或使用检查点恢复为止。 | 确保检查点存储位置可靠且可访问。 |
9.2 支持的检查点存储后端(Memory、File、Database)
| 存储类型 | 描述 | 示例代码 | 注意事项 |
|---|
| Memory | 使用内存作为临时存储,适用于测试环境。 | from langgraph.checkpoint.memory import MemoryCheckpoint
checkpoint = MemoryCheckpoint() | 数据仅存于当前运行实例,不适合生产环境。 |
| File | 将检查点保存至文件系统,适合小型应用。 | from langgraph.checkpoint.file import FileCheckpoint
checkpoint = FileCheckpoint(path="/tmp/checkpoint") | 需确保文件权限正确,避免并发写入冲突。 |
| Database | 利用数据库持久化检查点,支持大规模分布式部署。 | from langgraph.checkpoint.db import DBCheckpoint
checkpoint = DBCheckpoint(conn_str="postgresql://user:pass@localhost/db") | 配置适当的连接池大小和超时设置。 |
9.3 从检查点恢复执行
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 恢复执行 | app.resume(checkpoint_id) | 根据指定的检查点 ID 恢复流程。 | app.resume("checkpoint_12345") | 确认检查点存在且未损坏。 |
| 自动恢复 | 在构造函数中指定自动恢复策略 | 流程启动时自动尝试从最近的检查点恢复。 | app = workflow.compile(auto_resume=True, checkpoint=checkpoint) | 可能导致意外行为,需谨慎配置。 |
9.4 版本兼容与状态迁移
| 方法/概念 | 描述 | 注意事项 |
|---|
| 版本管理 | 管理不同版本间的检查点兼容性问题。 | 更新流程定义前应考虑现有检查点的影响。 |
| 状态迁移 | 当流程结构发生变化时,需要迁移旧检查点的数据格式。 | 实现迁移脚本或工具,逐步过渡老用户数据。 |
第10章:执行与调用(Invocation & Streaming)
10.1 同步执行:invoke() 方法
| 方法名 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
invoke | app.invoke(input_state) | 同步方式触发流程执行,并等待结果返回。 | result = app.invoke({"initial": "data"}) | 适用于快速响应场景,但可能阻塞主线程。 |
10.2 异步执行:ainvoke() 与异步支持
| 方法名 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
ainvoke | await app.ainvoke(input_state) | 异步方式触发流程执行,不阻塞主线程。 | async def main():
result = await app.ainvoke({"initial": "data"}) | 需要在异步环境中使用,如 async 函数内。 |
10.3 流式输出:stream() 与事件类型(on_chain_start, on_node_end 等)
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
stream | for event in app.stream(input_state): | 流式接收每个节点执行后的更新事件。 | for event in app.stream({"initial": "data"}):
print(event) | 适合实时监控和调试流程。 |
| 事件类型 | 包括但不限于 on_chain_start、on_node_end、on_error 等 | 提供详细的流程执行日志。 | if event["type"] == "on_node_end":
print(f"Node {event['node']} ended.") | 通过事件可以实现更精细的控制和反馈机制。 |
10.4 批量执行:batch() 与 abatch()
| 方法名 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
batch | app.batch([input_states]) | 同步批量处理多个输入状态。 | results = app.batch([{"id": 1}, {"id": 2}]) | 对于大量输入,考虑性能影响。 |
abatch | await app.abatch([input_states]) | 异步批量处理多个输入状态。 | async def main():
results = await app.abatch([{"id": 1}, {"id": 2}]) | 异步批量处理可提高效率,特别是在网络请求等 I/O 密集型任务中。 |
第11章:工具集成与 Agent 模式
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 工具定义 | 使用 @tool 装饰器或继承 BaseTool 类 | 定义可被调用的外部功能。 | from langchain_core.tools import tool
@tool
def search_web(query: str) -> str:
"""Search the web for information."""
return f"Results for {query}" | 工具需有明确输入输出类型,便于序列化和验证。 |
| 绑定工具至节点 | 将工具作为节点函数使用 | 在流程中直接调用工具。 | workflow.add_node("search", search_web) | 确保工具返回值符合状态更新格式(如字典)。 |
| 包装工具返回值 | 对工具结果进行封装后再更新状态 | 避免原始输出污染全局状态。 | def wrapped_search(state):
result = search_web.invoke(state["query"])
return {"tool_result": result} | 推荐做法,增强可控性和调试能力。 |
11.2 构建 ReAct 风格 Agent 流程
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| ReAct 循环结构 | Thought → Action → Observation → Repeat | 实现基于推理与行动的智能体循环。 | def decide_action(state):
if "need_info" in state:
return "search"
else:
return "respond"
workflow.add_conditional_edges("agent", decide_action, {"search": "search", "respond": END}) | 必须设置终止条件防止无限循环。 |
| 状态字段设计 | 维护 thoughts、actions、observations 字段 | 记录智能体决策过程。 | class State(TypedDict):
input: str
thoughts: list[str]
action: dict
observation: str
final_answer: str | 结构清晰有助于回溯和调试。 |
| 控制流组织 | 使用条件边驱动 ReAct 步骤跳转 | 明确各阶段执行顺序。 | workflow.add_edge("search", "observe")
workflow.add_edge("observe", "agent") # 回到决策点 | 观察后应回到决策节点继续判断下一步。 |
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| LLM 支持工具调用 | 使用支持 function calling 的模型(如 gpt-3.5-turbo) | 让 LLM 决定是否调用工具及参数。 | llm = ChatOpenAI(model="gpt-3.5-turbo").bind_tools([search_web])
def plan_with_llm(state):
msg = llm.invoke([HumanMessage(content=state["input"])])
return {"messages": [msg]} | 需绑定工具列表以便 LLM 可见可用功能。 |
| 解析工具调用请求 | 检查 LLM 输出中的 .tool_calls 属性 | 判断是否需要执行工具。 | if msg.tool_calls:
tool_call = msg.tool_calls[0]
result = search_web.invoke(tool_call)
return {"pending_tool_call": result} | 多工具场景下需遍历所有 tool_calls。 |
| 反馈观察结果给 LLM | 将工具执行结果以 AIMessage + ToolMessage 形式传回 | 支持多轮交互式推理。 | messages = [
HumanMessage(content="What's the weather in SF?"),
AIMessage(content="", tool_calls=[call]),
ToolMessage(content="Sunny", tool_call_id=call["id"])
]
final_response = llm.invoke(messages) | 消息历史必须正确传递才能维持上下文。 |
11.4 处理工具调用失败与重试
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 异常捕获 | try-except 包裹工具调用 | 防止因单个工具失败导致整个流程中断。 | def safe_tool_call(state):
try:
result = search_web.invoke(state["query"])
return {"result": result}
except Exception as e:
return {"error": str(e)} | 建议记录错误信息用于后续分析。 |
| 自动重试机制 | 使用 tenacity 或内置重试逻辑 | 提高容错性,适用于网络不稳定场景。 | from tenacity import retry, stop_after_attempt
@retry(stop=stop_after_attempt(3))
def reliable_search(query):
return search_web.invoke(query) | 设置最大尝试次数避免长时间等待。 |
| 错误边跳转 | 添加 error edge 到降级处理节点 | 实现失败后的优雅降级路径。 | workflow.add_error_edge("search", "fallback_handler") | 保证关键路径具备备选方案。 |
第12章:可观测性与调试
12.1 日志与回调系统(Callbacks)
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 回调注册 | 通过 config["callbacks"] 注册监听器 | 监听流程执行过程中的事件。 | from langchain.callbacks import StdOutCallbackHandler
app.invoke(input, config={"callbacks": [StdOutCallbackHandler()]}) | 支持多个回调处理器组合使用。 |
| 自定义回调类 | 继承 BaseCallbackHandler 并实现钩子方法 | 实现日志、指标收集等功能。 | class MyLogger(BaseCallbackHandler):
def on_chain_start(self, *args, **kwargs):
print("Chain started")
def on_node_end(self, *args, **kwargs):
print("Node ended") | 可用于审计、性能监控等高级用途。 |
12.2 集成 LangSmith 进行追踪
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| 启用 LangSmith 追踪 | 设置环境变量并启用 tracing | 自动上传执行轨迹到 LangSmith 平台。 | export LANGCHAIN_TRACING_V2=true
export LANGCHAIN_API_KEY=your_api_key | 需提前注册账号并获取 API Key。 |
| 手动配置 Tracer | 在调用时显式传入 tracer | 更细粒度控制追踪行为。 | from langchain.callbacks.tracers.langchain import wait_for_all_tracers
app.invoke(..., config={"callbacks": [...]})
wait_for_all_tracers() | 确保异步任务完成后再退出程序。 |
| 查看追踪详情 | 登录 https://smith.langchain.com | 分析每一步耗时、输入输出、错误等信息。 | — | 支持团队协作、测试比对、性能优化。 |
12.3 可视化图结构:graph.print_ascii() 与图形导出
| 方法/概念 | 语法 | 描述 | 示例代码 | 注意事项 |
|---|
| ASCII 图形输出 | graph.graph().print_ascii() | 在终端打印流程拓扑结构。 | workflow = StateGraph(State)
workflow.graph().print_ascii() | 适合快速查看结构,尤其在开发阶段。 |
| 导出为图像 | 使用 .draw_mermaid_png() 或 .draw_ascii() | 生成可视化图表用于文档或演示。 | graph_image = workflow.draw_mermaid_png()
with open("graph.png", "wb") as f:
f.write(graph_image) | 需安装额外依赖(如 pygraphviz、mermaid-cli)。 |
| Mermaid JS 兼容 | .draw_mermaid() 返回 Mermaid 语法字符串 | 可嵌入网页或 Markdown 文档。 | mermaid_code = workflow.draw_mermaid()
print(mermaid_code) | 支持在线预览(Mermaid Live Editor)。 |
12.4 调试常见错误:循环依赖、状态冲突、边未定义
| 错误类型 | 原因 | 解决方案 | 示例诊断 |
|---|
| 循环依赖 | 节点 A → B → C → A 无出口 | 添加最大步数限制或显式终止条件。 | MaxStepsExceeded 异常提示可能存在死循环。 |
| 状态冲突 | 多个节点同时修改同一字段导致不一致 | 使用不可变更新模式,避免直接修改原状态。 | 检查日志中状态变化顺序是否合理。 |
| 边未定义 | 条件函数返回了不存在的节点名 | 核对节点名称拼写,确保所有目标节点已添加。 | 报错:InvalidNextNode: Node 'unkown' not found(注意拼写错误)。 |
| 入口点未设置 | 忘记调用 set_entry_point() | 编译时报错提示缺少入口点。 | ValueError: Must set entry point before compiling. |
| 检查点未配置但启用中断 | 使用 interrupt_before 却未提供 checkpointer | 提供有效的检查点存储后端。 | ValueError: Cannot interrupt without a checkpointer. |
第13章:实战案例解析
13.1 构建多轮对话机器人
| 要素 | 描述 | 实现方式 | 注意事项 |
|---|
| 状态设计 | 维护 messages 列表与对话上下文字段 | 使用 MessageGraph 或自定义 State 包含 messages、user_intent、dialogue_stage | 确保消息历史不无限增长,可设置最大长度 |
| 节点划分 | 分为意图识别、回复生成、结束判断等节点 | workflow.add_node("classify_intent", classify_fn)
workflow.add_node("generate_response", llm_reply) | 意图识别可使用 LLM 或规则引擎 |
| 条件边控制流程 | 根据用户意图跳转不同处理分支 | def route_intent(state):
intent = detect_intent(state["messages"][-1].content)
return {"greeting": "respond", "query": "answer", "bye": "end"}[intent] | 需覆盖默认路径(如未知意图) |
| 终止条件 | 检测用户告别语或达到最大轮数 | workflow.add_conditional_edges("classify_intent", route_intent, ...)
workflow.add_edge("end_conversation", END) | 可结合 max_steps 防止无限对话 |
13.2 实现审批流程自动化
| 要点 | 描述 | 示例代码 | 注意事项 |
|---|
| 流程建模 | 模拟提交 → 审核 → 批准/拒绝 → 归档的完整流程 | class ApprovalState(TypedDict):
request: dict
approver: str
status: str # pending, approved, rejected
history: list | 明确每个状态的含义和转换规则 |
| 审批节点 | 人工审批环节支持中断等待 | app = workflow.compile(
checkpointer=MemoryCheckpoint(),
interrupt_before=["await_approval"]
) | 外部系统需能恢复执行(resume) |
| 多级审批 | 支持逐级审批逻辑 | def route_approver(state):
level = state["request"]["level"]
return f"approve_level_{level}" | 可结合角色权限系统动态路由 |
| 日志记录 | 记录每次审批操作的时间与人员 | def log_approval(state):
return {"history": [{"by": "alice", "action": "approved", "time": now()}]} | 保证审计可追溯性 |
13.3 构建数据分析助手(含工具调用链)
| 组件 | 描述 | 实现方式 | 注意事项 |
|---|
| 工具集定义 | 数据查询、可视化、解释等工具 | @tool
def query_db(sql: str) -> pd.DataFrame: ...
@tool
def plot_chart(data: dict) -> str: ... | 工具需处理异常并返回结构化结果 |
| ReAct 循环 | LLM 决策是否调用工具及参数 | llm = ChatOpenAI().bind_tools([query_db, plot_chart])
workflow.add_node("agent", lambda s: llm.invoke(s["messages"])) | 提示词应引导合理使用工具 |
| 工具执行节点 | 实际执行工具调用并将结果反馈给 LLM | def run_tool(state):
tool_call = state["messages"][-1].tool_calls[0]
result = tool_call.tool_fn(tool_call.args)
return {"messages": [ToolMessage(content=str(result), tool_call_id=tool_call.id)]} | 工具消息必须与 LLM 输出中的 tool_call_id 匹配 |
| 最终回答生成 | LLM 基于工具结果生成自然语言总结 | workflow.add_conditional_edges("agent", has_tool_call, {"yes": "run_tool", "no": "final_answer"}) | 避免工具调用后仍循环决策 |
13.4 分布式任务调度模拟
| 特性 | 描述 | 实现方式 | 注意事项 |
|---|
| 并行任务执行 | 同时启动多个独立任务 | workflow.add_parallel_nodes("tasks", ["task_a", "task_b", "task_c"], merge_results) | 任务间应无强依赖关系 |
| 子图封装 | 将每个任务封装为子图,便于复用 | task_graph = create_task_subgraph(task_id)
main_workflow.add_subgraph(task_graph, f"start_{task_id}") | 子图状态需与主图隔离或明确映射 |
| 失败重试机制 | 个别任务失败不影响整体流程 | workflow.add_error_edge("task_a", "retry_a") | 可配置最大重试次数和退避策略 |
| 结果聚合 | 所有任务完成后汇总结果 | def collect_results(states):
return {"results": [s["output"] for s in states]} | 聚合节点需等待所有并行路径完成 |
第14章:性能优化与最佳实践
14.1 状态精简与字段裁剪
| 方法 | 描述 | 示例 | 注意事项 |
|---|
| 移除冗余字段 | 不将临时或已处理的数据保留在状态中 | return {"processed_data": result, "temp_cache": None} | 减少序列化开销和存储压力 |
| 字段惰性加载 | 大数据字段按需加载而非全程携带 | def load_large_file(state):
if "file_ref" in state and "file_data" not in state:
return {"file_data": read_file(state["file_ref"])} | 适用于文件、图像等大对象 |
| 使用引用代替复制 | 状态中只保存 ID 或 URI | return {"report_id": "r_12345"} # 而非整个报告内容 | 外部服务负责数据存储与检索 |
14.2 缓存节点输出
| 方法 | 描述 | 实现方式 | 注意事项 |
|---|
| 内存缓存 | 使用 functools.lru_cache 缓存纯函数结果 | @lru_cache(maxsize=128)
def expensive_process(data): ... | 仅适用于无副作用、输入可哈希的函数 |
| 外部缓存 | 集成 Redis 或数据库缓存中间结果 | import redis
r = redis.Redis()
def cached_node(state):
key = hash_state(state)
if r.exists(key):
return r.get(key)
result = compute(state)
r.setex(key, 3600, serialize(result))
return result | 注意缓存失效策略和一致性 |
| 条件性缓存 | 根据输入特征决定是否缓存 | if state["mode"] == "dev":
return compute(state) # 不缓存开发模式 | 提高灵活性,避免污染测试数据 |
14.3 异步非阻塞设计
| 方法 | 描述 | 示例代码 | 注意事项 |
|---|
| 异步节点实现 | 使用 async def 定义 I/O 密集型节点 | async def fetch_data(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
return await resp.json() | 显著提升高并发场景下的吞吐量 |
| 异步工具调用 | 工具本身支持异步执行 | @tool
async def async_search(query): ... | 需 LLM 和运行环境支持 |
| 批量并发请求 | 使用 asyncio.gather 并行发起多个请求 | results = await asyncio.gather(
fetch_a(), fetch_b(), fetch_c()
) | 控制并发数防止资源耗尽(可用 Semaphore) |
14.4 图结构复用与模块化设计
| 方法 | 描述 | 示例 | 注意事项 |
|---|
| 子图封装通用流程 | 将验证、日志、通知等通用逻辑抽象为子图 | def create_validation_subgraph():
g = StateGraph(State)
g.add_node("validate", validate_fn)
g.add_conditional_edges("validate", ..., default="reject")
return g | 提高代码复用率,降低维护成本 |
| 参数化子图构建 | 通过参数定制子图行为 | def create_retry_loop(action_node, max_retries=3): ... | 增强灵活性,适应不同场景 |
| 图组合模式 | 主图通过条件边调用不同子图 | workflow.add_subgraph(validation_graph, "validate_input")
workflow.add_subgraph(approval_graph, "start_approval") | 保持主流程简洁清晰 |
| 版本化子图 | 对关键子图进行版本管理 | workflow.add_subgraph(v1_approval, "approval_v1")
workflow.add_subgraph(v2_approval, "approval_v2") | 支持灰度发布和回滚 |
第15章:扩展与生态集成
15.1 与 FastAPI 集成提供服务接口
将 LangGraph 工作流封装为 RESTful API,便于其他系统调用。
| 方法 | 描述 | 示例代码 | 注意事项 |
|---|
| 基础路由封装 | 使用 FastAPI 的 @app.post 暴露工作流入口 | from fastapi import FastAPI
from langgraph.graph import StateGraph
app = FastAPI()
@app.post("/chat")
async def chat_endpoint(input_data: dict):
result = workflow.invoke(input_data)
return result | 确保输入输出符合 JSON 序列化要求 |
| 异步支持 | 使用 ainvoke() 提高并发性能 | @app.post("/chat")
async def chat_endpoint(input_data: dict):
result = await workflow.ainvoke(input_data)
return result | 推荐用于生产环境,避免阻塞事件循环 |
| 请求校验 | 使用 Pydantic 模型验证输入 | from pydantic import BaseModel
class ChatInput(BaseModel):
message: str
session_id: str
@app.post("/chat")
async def chat_endpoint(data: ChatInput):
return await workflow.ainvoke({"input": data.message}) | 增强 API 安全性和健壮性 |
| 错误处理 | 捕获异常并返回标准错误响应 | from fastapi.exceptions import HTTPException
try:
result = await workflow.ainvoke(...)
except Exception as e:
raise HTTPException(status_code=500, detail=str(e)) | 避免泄露内部错误细节 |
15.2 与前端应用联动(WebSocket 流式响应)
通过 WebSocket 实现低延迟、双向通信的流式交互体验。
| 方法 | 描述 | 示例代码 | 注意事项 |
|---|
| WebSocket 路由 | 使用 FastAPI 的 @app.websocket 创建长连接 | from fastapi import WebSocket
@app.websocket("/ws/chat")
async def websocket_chat(websocket: WebSocket):
await websocket.accept()
async for event in workflow.astream(input_state):
await websocket.send_json(event) | 适合实时对话、进度反馈等场景 |
| 流式事件解析 | 从 stream() 输出中提取文本片段 | async for event in workflow.astream(...):
if "type" in event and event["type"] == "on_llm_new_token":
await websocket.send_text(event["content"]) | 可实现”打字机”效果 |
| 客户端 JavaScript | 前端监听并渲染流式消息 | const ws = new WebSocket("ws://localhost:8000/ws/chat");
ws.onmessage = (event) => {
const data = JSON.parse(event.data);
document.getElementById("output").innerText += data.content;
}; | 注意连接超时和重连机制 |
| 多会话隔离 | 使用 configurable 字段区分用户会话 | config = {"configurable": {"thread_id": session_id}}
async for event in app.astream(input, config): ... | 结合检查点实现会话持久化 |
15.3 自定义节点装饰器与中间件
通过装饰器和中间件增强节点功能,如日志、监控、权限等。
| 方法 | 描述 | 示例代码 | 注意事项 |
|---|
| 节点装饰器 | 为节点添加通用逻辑(如耗时统计) | import time
def timed_node(func):
async def wrapper(state):
start = time.time()
result = await func(state)
print(f"{func.__name__} took {time.time()-start:.2f}s")
return result
return wrapper
@timed_node
async def process_data(state): ... | 支持同步和异步函数包装 |
| 中间件模式 | 在执行前后插入逻辑(类似回调但更灵活) | def with_logging(node_fn):
def wrapped(state):
print(f"Entering {node_fn.__name__}")
result = node_fn(state)
print(f"Exiting {node_fn.__name__}")
return result
return wrapped | 可组合多个中间件 |
| 权限校验中间件 | 控制节点访问权限 | def require_role(role):
def decorator(node_fn):
def wrapped(state):
if state.get("user_role") != role:
raise PermissionError("Insufficient privileges")
return node_fn(state)
return wrapped
return decorator | 适用于多租户或分级系统 |
| 缓存中间件 | 自动缓存节点输出 | 见第14章缓存实践 | 需考虑缓存键生成策略 |
15.4 社区插件与第三方扩展
利用生态系统扩展 LangGraph 功能。
| 扩展类型 | 项目/平台 | 功能描述 | 集成方式 | 注意事项 |
|---|
| LangChain Tools | langchain-community | 数百种预构建工具(搜索、数据库、API等) | from langchain_community.tools import WikipediaQueryRun
from langchain_community.utilities import WikipediaAPIWrapper
wiki = WikipediaQueryRun(api_wrapper=WikipediaAPIWrapper()) | 检查依赖版本兼容性 |
| LangGraph Plugins | 社区贡献的图结构模板 | 如 react-agent-graph、hierarchical-agent 等 | 通过 pip 安装或直接导入模块 | 注意维护状态和安全性 |
| LangSmith 扩展 | smith.langchain.com | 提供追踪、评估、测试、监控能力 | 启用 LANGCHAIN_TRACING_V2 即可自动上报 | 支持自定义评估指标 |
| 可视化工具 | Mermaid.js、Draw.io、Diagrams.net | 将 .draw_mermaid() 输出可视化 | 复制 Mermaid 代码到编辑器预览 | 适合文档和演示 |
| 第三方服务集成 | Slack、Discord、Notion、Airtable | 将工作流接入常用协作平台 | 使用对应 SDK 或官方工具包 | 注意 API 调用频率限制 |