最近在做一个多步骤任务编排的需求,任务有明确的 SOP,但中间要等人审、要持久化进度、还有几段可以并行。一开始我想自己用循环撸,写了两版发现:流程控制不难,难的是中断恢复、检查点、并行汇聚这一套生产级能力。自己写一遍代价不低,而且多半会埋一堆坑。
于是去读了 spring-ai-alibaba-graph-core 这个模块的源码。它是 Spring AI Alibaba 整个框架的底层运行时引擎,负责图工作流编排和多智能体应用。我按「结构 → 执行流程 → 核心设计」的顺序过了一遍,重点想搞清楚三件事:状态怎么在节点间传、跑到一半怎么存档和恢复、怎么用响应式流把结果吐出去。这里把读到的内容整理下来。
它解决了什么问题
先说它到底提供了什么。对一个需要多步骤编排的 Java 应用,graph-core 给了四样东西:
- 有状态的工作流执行:支持中断、恢复、检查点保存
- 声明式图定义:通过节点(Node)和边(Edge)抽象复杂的多步骤任务流
- 持久化支持:内置 PostgreSQL、MySQL、Oracle、MongoDB、Redis、文件系统等存储后端
- 可观测性集成:与 Spring Cloud Observation / OpenTelemetry 无缝集成
这四点正好覆盖了我自己撸时会踩的那些坑。尤其第 1、3 点,是手写循环最容易出错、也最懒得补的地方。
顶层结构
先看目录结构,建立整体印象:
spring-ai-alibaba-graph-core/
├── src/main/java/com/alibaba/cloud/ai/graph/
│ ├── StateGraph.java # 核心:图定义入口
│ ├── CompiledGraph.java # 已编译的可执行图
│ ├── GraphRunner.java # 基于 Reactor 的执行引擎
│ ├── OverAllState.java # 全局状态容器
│ ├── RunnableConfig.java # 运行时配置
│ │
│ ├── action/ # 节点和边的执行动作接口
│ ├── checkpoint/ # 检查点系统(持久化)
│ │ └── savers/ # 多种存储实现
│ ├── executor/ # 图执行器(主图、子图)
│ ├── internal/ # 内部实现(节点、边的具体类)
│ ├── observation/ # 可观测性支持
│ ├── serializer/ # 状态序列化(Jackson、标准序列化)
│ ├── state/ # 状态管理策略
│ ├── store/ # 键值存储抽象
│ ├── skills/ # 技能注册与增强
│ └── scheduling/ # 定时 Agent 管理
│
└── src/test/java/ # 完善的测试套件(59+ 测试类)
几个关键入口对应着我关心的三个问题:
StateGraph.java:图定义 API,理解「怎么把流程写出来」GraphRunner.java:执行流程,理解「流程怎么跑起来」checkpoint/包:持久化机制,理解「跑到一半怎么存和恢复」
核心执行流程
执行流程遵循「定义 → 编译 → 执行」,底层用 Project Reactor 做响应式流处理。
graph TB
Start[定义 StateGraph] --> AddNodes[添加节点 addNode]
AddNodes --> AddEdges[添加边 addEdge/addConditionalEdges]
AddEdges --> Compile[编译 compile 生成 CompiledGraph]
Compile --> CreateRunner[创建 GraphRunner]
CreateRunner --> InitState[初始化 OverAllState]
InitState --> Execute[执行 run 方法]
Execute --> GraphRunnerContext[创建 GraphRunnerContext]
GraphRunnerContext --> MainExecutor[MainGraphExecutor.execute]
MainExecutor --> PreNode{节点前检查}
PreNode -->|中断| SaveCheckpoint[保存检查点]
PreNode -->|继续| ExecuteNode[执行节点 Action]
ExecuteNode --> UpdateState[更新 OverAllState]
UpdateState --> PostNode{节点后检查}
PostNode -->|中断| SaveCheckpoint
PostNode -->|继续| EvaluateEdge[评估边条件]
EvaluateEdge --> NextNode{下一个节点?}
NextNode -->|有| PreNode
NextNode -->|END| Return[返回 Flux GraphResponse]
SaveCheckpoint --> CheckpointSaver[(CheckpointSaver<br/>持久化层)]
CheckpointSaver --> Interrupt[中断并返回元数据]
style StateGraph fill:#e1f5ff
style CompiledGraph fill:#fff4e6
style GraphRunner fill:#f3e5f5
style CheckpointSaver fill:#e8f5e9
黄金路径分四步。我按自己跑通的最小例子来看:
1. 图定义阶段
StateGraph graph = new StateGraph(keyStrategyFactory)
.addNode("nodeA", nodeAction)
.addNode("nodeB", nodeAction)
.addEdge(START, "nodeA")
.addConditionalEdges("nodeA", routingFunction, Map.of(...))
.addEdge("nodeB", END);
2. 编译阶段
CompiledGraph compiledGraph = graph.compile(CompileConfig.builder()
.checkpointSaver(checkpointSaver) // 配置持久化
.build());
3. 执行阶段
GraphRunner runner = new GraphRunner(compiledGraph, runnableConfig);
Flux<GraphResponse<NodeOutput>> stream = runner.run(initialState);
4. 响应式流处理
GraphRunner委托给MainGraphExecutor- 执行器逐节点推进,每次发射
GraphResponse<NodeOutput> - 支持中断:遇到
InterruptableAction时保存检查点并暂停 - 支持恢复:使用保存的检查点 ID 恢复执行
把四步串起来看,定义是声明式的,执行是响应式的,中断和恢复是在节点前后两个检查点插进去的。这个结构清楚了,下面三个核心设计就好定位了。
三个核心设计
核心设计 1:状态管理与合并策略
位置:OverAllState.java、KeyStrategy.java、state/strategy/ 包
我自己写循环时第一个卡住的就是状态怎么传:直接覆盖会把上一步的结果丢了,自己累加又得处理并行节点谁先谁后。graph-core 的做法是采用键值对状态模型,每个键关联一个合并策略(KeyStrategy),控制状态在节点间传递时如何更新:
- ReplaceStrategy:直接替换(适用于单值状态)
- AppendStrategy:追加到集合(适用于消息历史)
- MergeStrategy:合并 Map(适用于复杂对象)
关键代码结构:
public final class OverAllState implements Serializable {
private Map<String, Object> data = new LinkedHashMap<>();
private Map<String, KeyStrategy> keyStrategies = new LinkedHashMap<>();
// 合并新值时应用策略
public <T> OverAllState mergeWith(String key, T value, KeyStrategy strategy) {
Object merged = strategy.apply(this.data.get(key), value);
this.data.put(key, merged);
return this;
}
}
把策略下放到每个键上,好处归结为三点:
- 避免覆盖冲突:尤其在并行节点中,多个分支同时写状态不会互相覆盖
- 支持累积模式:对话历史、审批链这类需要不断追加的场景能自然表达
- 状态可预测:合并规则是声明式注册的,不依赖执行顺序,跑出来的结果稳定
使用示例:
KeyStrategyFactory factory = KeyStrategyFactoryBuilder.builder()
.addDefault("chatHistory", NodeAggregationStrategy.APPEND) // 聊天历史追加
.addDefault("currentUser", NodeAggregationStrategy.REPLACE) // 用户信息替换
.build();
核心设计 2:检查点与中断恢复系统
位置:checkpoint/ 包、InterruptableAction.java
这是我最关心的部分,也是手写最难做对的地方。graph-core 支持可中断的工作流执行,关键组件有三个:
-
CheckpointSaver:抽象持久化接口
- 实现:
PostgresSaver、MysqlSaver、RedisSaver、MongoSaver、FileSystemSaver - 保存当前节点状态、下一节点 ID、完整状态快照
- 实现:
-
InterruptableAction:中断点接口
public interface InterruptableAction { // 节点执行前中断 Optional<InterruptionMetadata> interrupt(String nodeId, OverAllState state, RunnableConfig config); // 节点执行后中断 default Optional<InterruptionMetadata> interruptAfter(String nodeId, OverAllState state, Map<String, Object> actionResult, RunnableConfig config) { return Optional.empty(); } } -
恢复机制:
- 保存检查点时分配唯一 ID(UUID)
- 恢复时加载检查点,重建状态和节点位置
- 支持时间旅行:回溯到任意历史检查点
检查点数据结构:
public class Checkpoint {
private final String id; // 检查点 ID
private Map<String, Object> state; // 序列化的状态
private String nodeId; // 当前节点
private String nextNodeId; // 下一个节点
}
这套设计的价值,落到实际场景里就是三类需求:
- Human-in-the-Loop:在关键决策点暂停,等待人工审批
- 长时运行任务:支持任务跨会话持久化,服务重启不丢进度
- 调试与审计:可回溯任意执行点,定位问题不用靠猜
典型场景就是审批节点:
public class ApprovalNode implements AsyncNodeActionWithConfig, InterruptableAction {
@Override
public CompletableFuture<Map<String, Object>> apply(OverAllState state, RunnableConfig config) {
// 生成审批请求
return CompletableFuture.completedFuture(Map.of("request", "需要审批"));
}
@Override
public Optional<InterruptionMetadata> interruptAfter(...) {
// 执行后中断,等待人工审批
return Optional.of(InterruptionMetadata.builder(nodeId, state)
.addMetadata("reason", "待审批")
.build());
}
}
核心设计 3:响应式执行引擎(GraphRunner + MainGraphExecutor)
位置:GraphRunner.java、executor/MainGraphExecutor.java
第三个设计回答的是「结果怎么吐出去」。基于 Project Reactor 的响应式执行模型,支持:
- 流式输出:每个节点执行完立即发射结果
- 背压处理:下游消费速度控制上游生产
- 并行节点:多个独立节点同时执行,结果合并
执行器分层:
GraphRunner (外层 API)
↓
MainGraphExecutor (主图执行逻辑)
↓
NodeExecutor (单节点执行)
↓
ParallelEdgeProcessor (并行边处理)
核心执行循环:
public class MainGraphExecutor extends BaseGraphExecutor {
public Flux<GraphResponse<NodeOutput>> execute(GraphRunnerContext context, AtomicReference<Object> resultValue) {
return Flux.defer(() -> {
while (!isFinished(currentNodeId)) {
// 1. 检查中断点
if (shouldInterrupt(currentNodeId)) {
return interruptAndSave();
}
// 2. 执行节点
NodeOutput output = executeNode(currentNodeId, state);
// 3. 发射结果(流式输出)
emitResponse(GraphResponse.of(currentNodeId, output));
// 4. 更新状态
state = mergeNodeOutput(state, output);
// 5. 评估边,确定下一节点
currentNodeId = evaluateEdge(currentNodeId, state);
}
return Flux.empty();
});
}
}
响应式带来的好处,归结为三点:
- 实时反馈:前端可立即显示中间节点结果,不用等整条流程跑完
- 资源效率:避免一次性把所有结果加载到内存
- 灵活组合:可与 Spring WebFlux 无缝集成,直接做成 Server-Sent Events (SSE)
并行执行支持:
graph LR
A[nodeA] --> B[并行分支点]
B --> C[nodeB1]
B --> D[nodeB2]
B --> E[nodeB3]
C --> F[汇聚点]
D --> F
E --> F
F --> G[nodeC]
并行节点同时执行,结果通过状态合并策略汇聚——这里正好接上核心设计 1 的 KeyStrategy,并行写状态不冲突就是靠它兜住的。
架构全景图
把前面散落的组件拼到一起,整体长这样:
graph TB
subgraph "用户 API 层"
StateGraph[StateGraph<br/>图定义]
CompiledGraph[CompiledGraph<br/>已编译图]
GraphRunner[GraphRunner<br/>执行器]
end
subgraph "执行引擎层"
MainExecutor[MainGraphExecutor<br/>主执行器]
NodeExecutor[NodeExecutor<br/>节点执行器]
EdgeProcessor[ParallelEdgeProcessor<br/>边处理器]
end
subgraph "状态管理层"
OverAllState[OverAllState<br/>全局状态]
KeyStrategy[KeyStrategy<br/>合并策略]
StateSerializer[StateSerializer<br/>序列化器]
end
subgraph "持久化层"
CheckpointSaver[CheckpointSaver<br/>检查点保存器]
PostgresSaver[(PostgreSQL)]
MysqlSaver[(MySQL)]
RedisSaver[(Redis)]
MongoSaver[(MongoDB)]
FileSaver[(FileSystem)]
end
subgraph "可观测性层"
GraphObservation[GraphObservationHandler]
NodeObservation[NodeObservationHandler]
EdgeObservation[EdgeObservationHandler]
end
StateGraph --> CompiledGraph
CompiledGraph --> GraphRunner
GraphRunner --> MainExecutor
MainExecutor --> NodeExecutor
MainExecutor --> EdgeProcessor
NodeExecutor --> OverAllState
OverAllState --> KeyStrategy
OverAllState --> StateSerializer
MainExecutor --> CheckpointSaver
CheckpointSaver --> PostgresSaver
CheckpointSaver --> MysqlSaver
CheckpointSaver --> RedisSaver
CheckpointSaver --> MongoSaver
CheckpointSaver --> FileSaver
GraphRunner --> GraphObservation
NodeExecutor --> NodeObservation
EdgeProcessor --> EdgeObservation
style StateGraph fill:#e1f5ff
style CompiledGraph fill:#fff4e6
style GraphRunner fill:#f3e5f5
style OverAllState fill:#fce4ec
style CheckpointSaver fill:#e8f5e9
与替代方案的对比
读完之后横向比一下,看看 graph-core 在同类方案里的位置:
| 特性 | Graph Core | LangGraph (Python) | LangChain4j (Java) |
|---|---|---|---|
| 语言生态 | Java/Spring | Python | Java (原生) |
| 响应式流 | ✅ Reactor | ❌ 同步/asyncio | ❌ 同步 |
| 检查点持久化 | ✅ 多后端 | ✅ SQLite/Postgres | ⚠️ 基础支持 |
| Spring 集成 | ✅ 原生集成 | ❌ | ⚠️ 部分支持 |
| 并行节点 | ✅ 内置 | ✅ | ❌ |
| 可观测性 | ✅ Micrometer/OTel | ⚠️ LangSmith | ⚠️ 自定义 |
| Human-in-the-Loop | ✅ InterruptableAction | ✅ 断点机制 | ❌ |
回到我自己的场景,选 graph-core 的理由其实就几条:
- 已经是 Spring Boot 项目,技术栈一致
- 需要生产级的可观测性和监控
- 重视响应式流处理和实时反馈
- 需要企业级数据库持久化(PostgreSQL/Oracle)
总结
读完 graph-core,最深的体会是:一个工作流引擎的复杂度,不在流程控制本身,而在围绕流程的那套运行时能力。流程控制用循环也能写,但状态合并、检查点中断、响应式流式输出,这三块任何一个自己做都容易出 bug。
回头看,graph-core 三个核心设计对应的就是三个最难的运行时问题:
- KeyStrategy:状态怎么在节点(尤其并行节点)间安全传递
- CheckpointSaver + InterruptableAction:跑到一半怎么存档、怎么恢复、怎么等人审
- Reactor 执行引擎:结果怎么边算边吐、怎么背压、怎么并行汇聚
把这三块看明白,自己再写多步骤编排时,无论是直接用它,还是借鉴思路自己实现,心里都有底了。