LangGraph 异步引擎与检查点:打造可中断、可恢复的智能工作流

摘要
现代 AI 应用的工作流正变得愈加复杂:长时间调用外部 API、需要人工审核节点、任务中断后要能从上一次状态继续运行。LangGraph(基于 v0.2.30+ 版本验证)带来了全异步执行引擎和统一的检查点抽象,让开发者可以用几乎相同的方式构建“可暂停、可恢复、人机协同”的有状态智能体。本文通过一个“智能文章生成管道”实战,逐步拆解 async/await 节点、内存检查点、中断指令与动态子图,帮助你快速掌握其核心设计与最佳实践。


为什么需要异步和检查点?

如果你曾用 LangGraph 构建过生产级智能体,大概率遇到过这样几个痛点:

  • 耗时 I/O 阻塞调度:生成节点需要调用 OpenAI、Claude 等 API,同步执行会让整个图卡在该节点,浪费宝贵的计算资源。
  • 人工介入困难:某些节点需要人工确认(如“是否通过审核?”),传统做法是通过数据库轮询或外部回调来续跑,状态管理极易出错。
  • 断点续跑需求:调试时希望从失败的节点重新执行,而不是每次都从头跑一遍。

早期版本的 LangGraph 虽然提供了基础的 interrupt 和检查点接口,但开发者需要自行处理复杂的异步逻辑与线程安全。新版的关键变化在于:将异步执行与检查点深度集成到引擎内部,让整个图天然支持 async/await,并能以可插拔的方式持久化每一步的状态。


LangGraph 异步执行与检查点技术演进

当前版本(v0.2.30+)的执行引擎围绕两个核心抽象重构:

  1. 异步运行时
    图上的节点可以是 async def,执行器在调度时会自动 await,不阻塞异步事件循环。同时,同步节点仍可在异步图中运行,引擎负责统一调度。Command 机制也完全异步化,中断恢复后能够直接返回控制权到异步上下文。

  2. 检查点抽象
    所有状态都在每步执行前自动保存,包括当前节点 ID、状态快照和待处理队列。开发者只需指定一个 CheckpointSaver(支持内存、SQLite、Postgres 等),LangGraph 便会将所有通道的状态序列化到后端。当图因中断或异常暂停后,重新运行时只需传入相同的 thread_id 即可从断点处恢复,甚至支持“时间旅行”回退。

这种设计让你可以写出不受网络抖动、人工审核、长任务影响的健壮工作流,同时保持代码结构清晰。


实战:一个可中断的文章生成工作流

下面我们通过一个完整的例子来深入理解这些特性。工作流的需求如下:

  • 根据话题生成文章标题(异步调用外部 LLM 或搜索服务)。
  • 话题较短时走“纯文本”分支,较长时走“Markdown 带代码块”分支。
  • 生成摘要后,进入人工审核节点,支持人工确认或提供修改意见。
  • 审核通过后,由动态子图完成最终格式化并输出。

1. 全局状态与条件分支

首先定义贯穿整个流程的全局状态类型:

1
2
3
4
5
6
7
8
9
10
11
12
from typing import TypedDict, Literal

class ArticleState(TypedDict):
topic: str
classification: str # short / long
title: str
summary: str
format_type: str # markdown / plain
need_code_block: bool # 子图使用
formatted: str
# 以下用于人工审核中断
human_feedback: str

条件分类节点根据话题长度决定后续路径,这是一个同步节点,但可以无缝嵌入异步图:

1
2
3
4
5
6
7
8
9
10
11
12
def classify_topic(state: ArticleState) -> dict:
if len(state["topic"]) < 20:
return {
"classification": "short",
"format_type": "plain",
"need_code_block": False
}
return {
"classification": "long",
"format_type": "markdown",
"need_code_block": True
}

2. 异步节点与外部模拟

生成标题的节点模拟了一次异步 API 调用。在当前版本中,你只需要将节点定义为 async def

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
import asyncio

async def generate_title(state: ArticleState) -> dict:
# 模拟外部 API 延迟,如调用 OpenAI
await asyncio.sleep(0.5)
topic = state["topic"]
if state["classification"] == "long":
title = f"深入解析:{topic} 的核心原理与实践"
else:
title = f"快讯:{topic} 指南"
return {"title": title}

async def generate_summary(state: ArticleState) -> dict:
await asyncio.sleep(0.5)
summary = f"本文围绕「{state['title']}」展开讨论,分析了关键架构与最佳实践。"
return {"summary": summary}

引擎在执行这两个节点时会自动挂起协程,不会阻塞其他并行任务(如果有)。

3. 子图:动态格式化引擎

为了演示复杂编排,我们将格式化逻辑封装成一个独立的子图。子图状态定义与全局主图略有不同,只暴露必要字段:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
# subgraphs.py
from typing import TypedDict
from langgraph.graph import StateGraph, END

class FormatState(TypedDict):
title: str
summary: str
need_code_block: bool
output: str

def _markdown_formatter(state: FormatState) -> dict:
title = state["title"]
summary = state["summary"]
code_block = ""
if state["need_code_block"]:
code_block = (
"\n```python\n"
"# 示例代码\n"
"from langgraph import StateGraph\n"
"```\n"
)
output = f"## {title}\n\n{summary}{code_block}"
return {"output": output}

def _plain_formatter(state: FormatState) -> dict:
output = f"{state['title']}\n{state['summary']}"
return {"output": output}

def build_formatting_subgraph(fmt_type: str) -> StateGraph:
"""根据格式类型动态构建子图,返回编译好的 Graph 实例"""
builder = StateGraph(FormatState)
if fmt_type == "markdown":
builder.add_node("formatter", _markdown_formatter)
else:
builder.add_node("formatter", _plain_formatter)
builder.set_entry_point("formatter")
builder.add_edge("formatter", END)
return builder.compile()

子图本身可以作为主图的一个节点——LangGraph 原生支持这种嵌套,让主图只需通过一个 add_node 即可将整个子流程抽象为一个黑盒。

4. 中断与人工审核节点

这是本文最亮眼的部分。LangGraph 将 interrupt 提升为一等公民,它可以在任意节点中暂停执行,并在恢复时通过 Command 注入外部数据。

在标题和摘要生成之后,我们插入一个人工审核节点:

1
2
3
4
5
6
7
8
9
10
11
12
13
from langgraph.types import interrupt, Command

def human_review(state: ArticleState) -> dict:
# 中断执行,将当前状态通过中断值抛给外部(如前端UI)
# interrupt() 会挂起当前协程,直到外部通过 Command(resume=...) 恢复
# 恢复后,外部传入的值将作为 interrupt() 的返回值赋给 feedback
feedback = interrupt({
"title": state["title"],
"summary": state["summary"],
"message": "请审核文章内容,输入 'approve' 通过或提供修改意见"
})
# 在外部恢复后,此节点才会继续执行,将反馈写入状态并返回
return {"human_feedback": feedback}

当执行到 interrupt() 时,当前协程会完全挂起,整个图的状态会被自动保存到检查点。此时,该节点实际上尚未执行 return 语句;外部系统可以获取中断信号,并等待决策。当调用 graph.ainvoke(Command(resume="在标题中增加'2025版'"), config) 时,interrupt() 会解除阻塞,返回字符串 "在标题中增加'2025版'",然后节点继续运行,将 feedback 的值写入状态并返回 {"human_feedback": feedback}。这种机制确保了状态更新的时机清晰、数据流可控。

5. 构建主图并注入检查点

最后,我们拼装主图,并指定内存检查点存储器 MemorySaver(该类原生支持异步环境,无需区分 Async 前缀):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
import asyncio
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from subgraphs import build_formatting_subgraph

# 构建主图
builder = StateGraph(ArticleState)

builder.add_node("classify", classify_topic)
builder.add_node("generate_title", generate_title)
builder.add_node("generate_summary", generate_summary)
builder.add_node("human_review", human_review)
# 子图作为节点:这里简单地选择一种格式化方式,实际可根据 format_type 动态添加
builder.add_node("format_output", build_formatting_subgraph("markdown"))

# 连线
builder.add_edge(START, "classify")
builder.add_edge("classify", "generate_title")
builder.add_edge("generate_title", "generate_summary")
builder.add_edge("generate_summary", "human_review")
builder.add_conditional_edges(
"human_review",
lambda s: "format_output" if s["human_feedback"] == "approve" else END,
{"format_output": "format_output", END: END}
)
builder.add_edge("format_output", END)

# 使用内存检查点存储器(同步/异步通用)
memory = MemorySaver()
graph = builder.compile(checkpointer=memory)

注意 compilecheckpointer 参数:只要传入一个 CheckpointSaver 实例,LangGraph 就会在每步执行前后自动保存检查点。恢复时只需使用相同的 thread_id

1
config = {"configurable": {"thread_id": "article-1"}}

运行效果演示

启动任务后,首次调用会一直执行到 human_review 节点并中断。由于 interrupt() 会引发 GraphInterrupt 异常,我们需要捕获该异常并获取当前状态,随后通过 Command(resume=...) 恢复执行。以下是符合最佳实践的示例代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
from langgraph.errors import GraphInterrupt
from langgraph.types import Command

async def main():
config = {"configurable": {"thread_id": "article-1"}}

# 首次执行,预期会在 human_review 节点中断
try:
await graph.ainvoke({"topic": "LangGraph 异步引擎设计"}, config)
except GraphInterrupt:
# 中断后,可获取检查点中保存的最新状态
state_snapshot = await graph.aget_state(config)
print("中断时状态:", state_snapshot.values)

# 模拟人工审核通过,恢复执行
await graph.ainvoke(Command(resume="approve"), config)
final_state = await graph.aget_state(config)
print("最终输出:", final_state.values.get("formatted"))

控制台输出类似:

1
2
3
4
5
6
7
8
9
中断时状态: {'topic': 'LangGraph 异步引擎设计', 'classification': 'long', 
'title': '深入解析:LangGraph 异步引擎设计 的核心原理与实践',
'summary': '本文围绕「深入解析:...」展开讨论...', 'human_feedback': None}
最终输出: ## 深入解析:LangGraph 异步引擎设计 的核心原理与实践

本文围绕...讨论。
```python
# 示例代码
from langgraph import StateGraph

- 如果人工输入的是修改意见(如 `"修改标题"`),你还可以在后续节点中根据 `human_feedback` 动态调整状态。
- 因为检查点的存在,哪怕程序崩溃或手动停止,重新执行时只要传入同一个 `thread_id`,工作流会从上次中断处继续,而不是重新开始。

---

## 总结与展望

LangGraph 的异步引擎与检查点抽象,将原本需要大量胶水代码才能完成的“可恢复、可暂停、可协作”智能体模式,变成了引擎的原生能力。核心得益点在于:

- **统一异步模型**:任何节点都可以 `async`,且同步节点自动兼容,降低了迁移成本。
- **标准化中断与恢复**:`interrupt` + `Command` 让外部交互变得干净自然,不再依赖轮询或手工状态拼接。
- **可插拔检查点**:内存、SQLite、Postgres 等存储器对业务代码完全透明,只需修改一行配置即可实现持久化,为生产环境提供了坚实保障。

未来版本可能会带来更多高级特性,例如更细粒度的检查点策略、并行分支中断以及分布式检查点后端等。但就目前而言,新版已经为构建稳定、健壮的 LLM 工作流打下了绝佳基础。如果你还在使用旧版同步执行模式,现在正是升级的最佳时机。