graph-core 的中断和检查点机制解决的是 AI 工作流里的两个核心问题:
- 高风险操作需人工确认:执行 SQL、调用工具、支付转账这类操作,执行前要能暂停等人审
- 长流程崩溃后要能从断点恢复:服务重启后不能从头重跑,否则前面 LLM 调用的成本白花、副作用操作可能重复
源码里用一套三层机制把它们一起解决:
- 中断(Interrupt):在执行链上制造暂停点
- 检查点(Checkpoint):把暂停那一刻的状态落盘
- 恢复(Resume):从落盘状态接着跑
这里把读到的内容整理下来,重点在三个中断时机的边评估时机差异,也是最容易用错的地方。
三层机制怎么协作
先给整体结构,三层各管一件事:
第 1 层:中断机制(Interrupt)
在执行链上制造可恢复的暂停点;提供 Before / After / BeforeEdge 三种时机
第 2 层:检查点持久化(Checkpoint)
保存中断时的完整状态快照;存储可插拔(内存 / 数据库)
第 3 层:恢复机制(Resume)
从检查点重建执行上下文并继续;支持状态修改 + 手动触发
三层是配合用的:中断制造暂停点,检查点把暂停那一刻的状态落盘,恢复从落盘的状态接着跑。分开看每一层都不复杂,组合起来才能支撑「长生命周期、可暂停、可审阅、可修改、可恢复」的工作流。
第一层:中断的三个时机
这是整个设计里最容易绕进去的地方。graph-core 提供三种中断时机:
CompileConfig.builder()
.interruptBefore("dangerousNode") // 执行前中断
.interruptAfter("reviewNode") // 执行后中断
.interruptBeforeEdge(true) // 边路由前中断
.build();
三种时机对应三种诉求:
| 中断类型 | action 执行? | 状态合并? | 边评估? | nextNodeId | 适用场景 |
|---|---|---|---|---|---|
interruptBefore("C") |
未执行 | 否 | 是 | "C" |
阻止高风险操作执行 |
interruptAfter("B") |
已执行 | 是 | 是 | 已确定 | 审查执行结果 |
interruptAfter("B") + interruptBeforeEdge(true) |
已执行 | 是 | 否 | null |
干预下一步路由决策 |
节点 interruptAfter() |
已执行 | 是 | 是 | 已确定 | 节点根据结果动态中断 |
前两种好理解:Before 在节点执行前停,能阻止危险动作;After 在节点执行后停,能看到结果再决定。BeforeEdge 是第三种,也是最绕的一个——它修饰的是 interruptAfter,意思是「节点执行完了,但还没评估边、还没定下一跳」。这个时机的价值在于:用户可以趁机改状态,从而改变路由走向。
除了配置触发,还支持节点代码里手动标记,适合动态判断:
.addNode("conditionalNode", (state, config) -> {
if (isHighRisk(state)) {
config.markNodeAsInterrupted("conditionalNode"); // 运行时决定中断
}
return Map.of("result", "...");
})
配置触发适合固定的审批点,手动触发适合「风险评分超阈值才暂停」这类动态判断。两种方式各有适用面,同时提供。
第二层:检查点每个节点都存
检查点的核心决策是「何时存」。graph-core 的选择是每个节点执行后自动存一份:
// NodeExecutor 执行流程
1. 执行节点 action
2. 合并状态 mergeIntoCurrentState(updateState)
3. 计算下一节点 nextNodeId
4. 保存 Checkpoint(自动) // 关键:自动保存
5. 返回 NodeOutput
检查点里存了什么:
public class Checkpoint {
private final String id; // UUID(每次保存生成新 ID)
private Map<String, Object> state; // 完整状态快照
private String nodeId; // 当前节点 ID
private String nextNodeId; // 下一个节点 ID
}
为什么每个节点都存,而不是只在中断时存?对比一下三种策略:
| 策略 | 优势 | 劣势 | graph-core 选择 |
|---|---|---|---|
| 只在中断时保存 | 节省存储 | 崩溃时丢失未保存的进度 | 否 |
| 每个节点都保存 | 任意节点崩溃都能恢复 | 存储开销大 | 是 |
| 用户手动保存 | 最灵活 | 用户容易忘记 | 否 |
这个选择的底色是「可靠性优先于性能」。AI 工作流单次 LLM 调用就要花几分钱到几块钱,长流程可能几十次调用,崩溃重跑的成本远高于存储成本。再加上保留全部历史检查点,还顺带支持了「时间旅行」——回溯到任意历史节点做调试或分叉。
存储位置走可插拔策略,内置七种实现:
BaseCheckpointSaver (接口)
├── MemorySaver // 内存(仅测试/调试)
├── FileSystemSaver // 本地文件
├── PostgresSaver // PostgreSQL
├── MysqlSaver // MySQL
├── OracleSaver // Oracle
├── MongoSaver // MongoDB
└── RedisSaver // Redis
| 存储 | 持久化 | 跨进程 | 性能 | 适用场景 |
|---|---|---|---|---|
| MemorySaver | 否 | 否 | 极快 | 测试、调试 |
| FileSystemSaver | 是 | 受限 | 快 | 单机部署 |
| RedisSaver | 是 | 是 | 极快 | 高性能、短期状态 |
| PostgresSaver | 是 | 是 | 中 | 生产环境、长期归档 |
接口定义统一,切换存储不用改业务代码:
public interface BaseCheckpointSaver {
Collection<Checkpoint> list(RunnableConfig config); // 列出所有 Checkpoint
Optional<Checkpoint> get(RunnableConfig config); // 获取当前 Checkpoint
RunnableConfig put(RunnableConfig config, Checkpoint cp); // 保存 Checkpoint
Tag release(RunnableConfig config); // 释放资源
}
每个节点都存的代价是存储会涨,所以生命周期可配置:
RunnableConfig.builder()
.metadata(CHECKPOINTS_NUM_RETAINED, 10) // 只保留最近 10 个
.build();
// 框架按配置清理最旧的
default void retainLatestCheckpoints(LinkedList<Checkpoint> checkpoints, RunnableConfig config) {
checkpointsNumRetained(config).ifPresent(numRetained -> {
while (checkpoints.size() > numRetained) {
checkpoints.removeLast();
}
});
}
不同场景的保留策略可以差很多:调试/审计要全留,生产对话留最近 20 个,纯崩溃恢复的批处理留最新一个就够。
第三层:恢复靠手动触发
恢复的设计要点是两个:谁触发、能不能改状态。
第一个问题,谁触发。graph-core 选的是手动触发,用户主动调用:
// 第一阶段:正常执行,遇到中断点暂停
workflow.stream(initialState, config).collectList().block();
// 返回 InterruptionMetadata,流暂停
// 用户审查、修改状态后,手动恢复
RunnableConfig resumeConfig = RunnableConfig.builder().resume().build();
workflow.stream(null, resumeConfig) // 手动触发恢复
.collectList().block();
// 从中断点继续执行
为什么不自动恢复?因为 Human-in-the-Loop 的核心就是人参与决策。中断一暂停就自动继续,审查这一步就形同虚设。手动触发让用户掌控「何时恢复」和「是否恢复」。
第二个问题,能不能改状态。答案是能,这正是这套机制区别于「单纯暂停」的地方:
// CompiledGraph.java:309-332
public RunnableConfig updateState(RunnableConfig config,
Map<String, Object> values,
String asNode) throws Exception {
// 1. 获取当前 Checkpoint
Checkpoint checkpoint = saver.get(config).map(Checkpoint::copyOf).orElseThrow();
// 2. 合并用户修改的值(使用 KeyStrategy)
checkpoint = checkpoint.updateState(values, keyStrategyMap);
// 3. 可选:指定下一个节点(覆盖原有路由)
if (asNode != null) {
Command nextCommand = nextNodeId(asNode, checkpoint.getState(), config);
checkpoint = checkpoint.updateState(nextCommand.update(), keyStrategyMap);
}
// 4. 保存修改后的 Checkpoint
return saver.put(config, checkpoint);
}
支持的修改操作:
| 修改类型 | API | 用途 |
|---|---|---|
| 修改状态值 | updateState(config, Map.of("key", "newValue")) |
纠正错误、补充信息 |
| 改变路由 | updateState(config, values, "nextNode") |
强制跳转到其他节点 |
| 修改后恢复 | updateState() + resume() |
应用修改并继续执行 |
落到那个 SQL 删除的场景上,就是中断后把误生成的 DELETE FROM users WHERE 1=1 改成 DELETE FROM users WHERE id = 123,再恢复执行:
// 中断时的状态
// state: {"sql": "DELETE FROM users WHERE 1=1"} // 错误
// 用户修正
config = workflow.updateState(config, Map.of(
"sql", "DELETE FROM users WHERE id = 123"
));
// Resume 继续执行(基于修正后的状态)
workflow.stream(null, config.withResume());
中断时返回什么
中断触发后返回的是 InterruptionMetadata,不是单纯一个节点 ID:
public final class InterruptionMetadata extends NodeOutput {
private final Map<String, Object> metadata; // 自定义元数据
private List<ToolCall> toolsAutomaticallyApproved; // 工具调用列表
private List<ToolFeedback> toolFeedbacks; // 工具反馈
}
包含的信息:
- 当前节点 ID:
nodeId(),中断发生在哪个节点 - 当前状态:
state(),完整的 OverAllState 快照 - 元数据:
metadata(),自定义信息(如原因、建议) - 工具调用:
toolsAutomaticallyApproved,即将执行的工具列表
用户做审批决策需要完整上下文,光知道「停在某节点」不够,得看到「即将执行什么 SQL」「为什么被判定为高风险」。元数据字段允许节点附带自定义提示。
完整生命周期
把三层机制串起来,中断流程的完整时序是这样的:
sequenceDiagram
participant User
participant GraphRunner
participant NodeExecutor
participant CheckpointSaver
User->>GraphRunner: invoke(initialState)
GraphRunner->>NodeExecutor: execute(NodeA)
NodeExecutor->>CheckpointSaver: put(checkpoint_A)
NodeExecutor-->>GraphRunner: NodeOutput
GraphRunner->>NodeExecutor: execute(NodeB)
Note over GraphRunner: 检测到 interruptAfter("NodeB")
NodeExecutor->>CheckpointSaver: put(checkpoint_B)
GraphRunner-->>User: 返回 InterruptionMetadata
Note over User: 用户审查状态<br/>决定:批准/拒绝/修改
alt 用户批准
User->>GraphRunner: stream(null, resume=true)
GraphRunner->>CheckpointSaver: get(checkpoint_B)
GraphRunner->>NodeExecutor: execute(NodeC)
NodeExecutor-->>User: 最终结果
else 用户修改
User->>GraphRunner: updateState(修改值)
GraphRunner->>CheckpointSaver: put(checkpoint_B_modified)
User->>GraphRunner: stream(null, resume=true)
GraphRunner->>NodeExecutor: execute(NodeC) [基于修改后的状态]
else 用户拒绝
Note over User: 不调用 resume,流程终止
end
一个关键分工是:中断和检查点保存是自动的,恢复和状态修改是手动的。框架负责可靠性(自动保存),用户负责决策(审查后恢复)。
触发主体和时机的对照如下:
| 操作 | 触发方 | 触发时机 | 触发方式 |
|---|---|---|---|
| 配置中断点 | 用户 | 编译时 | CompileConfig.interruptAfter("node") |
| 检测中断 | Graph 引擎 | 每个节点执行前/后 | context.shouldInterrupt() |
| 触发中断 | Graph 引擎 | 检测到中断点 | 返回 InterruptionMetadata |
| 保存检查点 | Graph 引擎 | 每个节点执行后 | 自动调用 saver.put() |
| 修改状态 | 用户 | 中断后、Resume 前 | updateState(config, values) |
| 恢复 | 用户 | 中断后 | stream(null, resume=true) |
三个 BeforeEdge 误用的坑
三种中断时机里,interruptBeforeEdge 我自己用错过两次,单独拎出来讲。它有几个反直觉的点:
坑 1:以为 Resume 会重新执行中断节点的 action
错误理解是「B 中断 → Resume 时重新执行 B.action」。实际流程是「B 中断 → Resume 从 checkpoint.nextNodeId 继续,不重新执行 B」。interruptAfter 的语义是「节点已执行完」,不是「节点待执行」。
坑 2:以为 interruptBeforeEdge(true) 能单独用
.interruptBeforeEdge(true) // 单独使用无意义
interruptBeforeEdge 只是一个修饰符,控制「如果中断,是在边评估前还是后」。没有 interruptAfter 触发中断,这个开关不生效。正确写法是配合 interruptAfter:
.interruptAfter("B") // 必须配合使用
.interruptBeforeEdge(true) // 修饰 interruptAfter 行为
坑 3:以为改状态总能改变路由
最隐蔽的一个。期望「中断后 updateState 改了路由字段,Resume 就会走新分支」,但实际只有配了 interruptBeforeEdge(true) 时才成立:
| 配置 | 修改状态能改变路由? | 原因 |
|---|---|---|
interruptAfter("B") |
否 | nextNodeId 已固化 |
interruptAfter("B") + interruptBeforeEdge(true) |
是 | nextNodeId 为 null,Resume 时重新评估 |
interruptBefore("C") |
否 | nextNodeId 已固化为 “C” |
原因是边评估的时机不同。普通 interruptAfter 在节点执行后就把边评估了,nextNodeId 已写死进检查点;interruptBeforeEdge(true) 把边评估推迟到 Resume 之后,Resume 时重新读状态算路由,改状态才有意义。
对照场景看更清楚——节点 B 根据 messages 字段选走 C 或 D:
配置 interruptAfter("B")(默认):
B 执行 → 合并状态(messages="B") → 评估边 → nextNodeId="D" → 中断
Resume → 直接执行 D
updateState(messages="C") 无法改变路由(D 已固化)
配置 interruptAfter("B") + interruptBeforeEdge(true):
B 执行 → 合并状态(messages="B") → 中断(边未评估)
Resume → 重新评估边 → 从 messages 读取
updateState(messages="C") 可以改变路由(走 C)
理解了边评估时机,三个时机的关系就清楚了:
Before = 还没干(节点未执行)
After = 已经干完(节点执行完、边也评估完)
BeforeEdge = 干完了但路还没选(节点执行完、边未评估)
检查点字段怎么被 Resume 用
理解了上面的坑,再看检查点的 nodeId 和 nextNodeId 两个字段就顺了。关键是 nodeId 不等于「中断发生的节点」,而是「最后一个已完成的节点」;nextNodeId 是「准备执行但还没执行的节点」。
不同中断位置下检查点内容不同:
| 中断位置 | checkpoint.nodeId | checkpoint.nextNodeId | checkpoint.state |
|---|---|---|---|
interruptBefore("C") |
"B" |
"C" |
B 执行后的状态 |
interruptAfter("B") |
"B" |
"D"(已评估) |
B 执行后的状态 |
interruptAfter("B") + beforeEdge |
"B" |
null(未评估) |
B 执行后的状态 |
Resume 的源码逻辑(GraphRunnerContext.java:96-118)就是按这两个字段重建上下文:
checkpoint = saver.get(config)
currentNodeId = null // 重置
nextNodeId = checkpoint.getNextNodeId() // 从 checkpoint 恢复
overallState = initialState.input(checkpoint.getState()) // 恢复状态
流程是:从 saver 取检查点 → 恢复 nextNodeId(为 null 时会重新评估边)→ 恢复状态 → 交给执行器继续执行 nextNodeId。
和传统工作流的对比
读完横向比一下,看看这套设计和传统工作流引擎的差别:
| 特性 | 传统工作流(Camunda/Activiti) | graph-core |
|---|---|---|
| 中断机制 | 「用户任务」(User Task) | Before/After/BeforeEdge 三种时机 |
| 状态持久化 | 流程实例表 | 可插拔 Saver(7 种存储) |
| 状态修改 | 有限支持(需要管理 API) | 原生支持 updateState() |
| 时间旅行 | 不支持 | 保留所有 Checkpoint 可回溯 |
| 编程模型 | BPMN XML | 代码定义(类型安全) |
| AI 集成 | 需要自己实现 | 原生支持(设计目标就是 AI 工作流) |
graph-core 的设计目标是 AI 工作流,所以重点放在了 LLM 调用昂贵、工具调用有副作用、状态需要可审计这些诉求上。传统工作流引擎在管理界面和企业集成上更成熟,但状态修改和时间旅行这类能力是缺的。
graph-core 也借鉴了 LangGraph 的思路,但在 Java 生态里做了适配:
| 特性 | LangGraph | graph-core |
|---|---|---|
| 中断时机 | 只有 interrupt_before/after |
额外支持 interruptBeforeEdge |
| 状态修改 | 通过 Python dict | 通过 KeyStrategy 保证合并语义 |
| 类型安全 | Python 动态类型 | Java 静态类型 |
| 持久化 | 主要 PostgreSQL | 支持 7 种存储 |
落到实际场景
这套机制适用的场景,回到我自己的几个需求上:
代码审查 Agent:AI 生成代码后中断等人审,审查通过后自动跑测试,部署前再次中断确认。
StateGraph codeReviewWorkflow = new StateGraph(factory)
.addNode("generateCode", aiGenerateCode)
.addNode("runTests", executeTests)
.addNode("deploy", deployToProduction)
.addEdge(START, "generateCode")
.addEdge("generateCode", "runTests")
.addEdge("runTests", "deploy")
.addEdge("deploy", END)
.compile(CompileConfig.builder()
.interruptAfter("generateCode") // 代码生成后人工审查
.interruptBefore("deploy") // 部署前二次确认
.build());
金融风控审批:低风险自动通过,高风险动态触发人工审批,交易前强制确认。
StateGraph riskWorkflow = new StateGraph(factory)
.addNode("calculateRisk", aiCalculateRisk)
.addNode("approveOrReject", (state, config) -> {
double riskScore = state.value("riskScore");
if (riskScore > 80) {
config.markNodeAsInterrupted("approveOrReject"); // 高风险需人工
}
return Map.of("autoDecision", riskScore < 30 ? "approve" : "pending");
})
.addNode("executeTransaction", executeTransaction)
.compile(CompileConfig.builder()
.interruptBefore("executeTransaction") // 交易前最后确认
.build());
不适用的情况也要说清:纯计算无副作用的流程、毫秒级实时响应、完全自动化无需人工,这些没必要配中断点,硬加只会增加延迟和复杂度。
一个最小可用模式
把三层机制抽出来,核心模式可以浓缩成下面这段伪代码,自己实现工作流时可以直接借鉴:
class CheckpointedWorkflow {
CheckpointSaver saver;
Set<String> interruptPoints;
void execute(State initialState, String sessionId) {
State current = initialState;
String nodeId = START;
while (nodeId != END) {
State newState = executeNode(nodeId, current);
saver.save(sessionId, nodeId, newState); // 每个节点自动保存
if (interruptPoints.contains(nodeId)) {
throw new InterruptException(nodeId, newState); // 暂停并返回
}
nodeId = routeNext(nodeId, newState);
current = newState;
}
}
void resume(String sessionId, State modifications) {
Checkpoint cp = saver.load(sessionId);
State modifiedState = cp.state.merge(modifications); // 应用用户修改
execute(modifiedState, cp.nextNodeId, sessionId);
}
}
四个关键要素:每个节点自动存检查点、编译时声明中断点、Resume 前可改状态、按 sessionId 隔离不同执行实例。
总结
中断、检查点、恢复三层机制共同构成了 graph-core 的 Human-in-the-Loop 能力,各自职责清晰:
- 中断:在执行链上制造暂停点,Before / After / BeforeEdge 三种时机分别覆盖「阻止执行 / 审查结果 / 干预路由」三类诉求。其中
BeforeEdge的设计最为关键,它将边评估推迟到 Resume 之后,使用户修改状态能实际影响路由决策,这是普通interruptAfter做不到的。 - 检查点:采用「每节点自动保存」策略,以存储成本换取可靠性,任意节点崩溃均可恢复至最近状态,同时保留全部历史以支持时间旅行与审计回溯。存储层可插拔,生命周期可配置,适配从调试到生产的不同场景。
- 恢复:手动触发,恢复前支持状态修改与路由覆盖。自动保存与手动恢复的分工,将可靠性交给框架、决策权交给用户,符合 Human-in-the-Loop 的本质。
这套设计对 AI 工作流的价值在于:它把「流程可控」与「状态可靠」两个通常割裂的工程问题统一到了一套机制里。理解三个中断时机的边评估时机差异,是正确使用这套机制的关键,也是最容易出错的地方。