最近在做一个多步骤任务编排的需求,任务有明确的 SOP,但中间要等人审、要持久化进度、还有几段可以并行。一开始我想自己用循环撸,写了两版发现:流程控制不难,难的是中断恢复、检查点、并行汇聚这一套生产级能力。自己写一遍代价不低,而且多半会埋一堆坑。

于是去读了 spring-ai-alibaba-graph-core 这个模块的源码。它是 Spring AI Alibaba 整个框架的底层运行时引擎,负责图工作流编排和多智能体应用。我按「结构 → 执行流程 → 核心设计」的顺序过了一遍,重点想搞清楚三件事:状态怎么在节点间传、跑到一半怎么存档和恢复、怎么用响应式流把结果吐出去。这里把读到的内容整理下来。

它解决了什么问题

先说它到底提供了什么。对一个需要多步骤编排的 Java 应用,graph-core 给了四样东西:

  1. 有状态的工作流执行:支持中断、恢复、检查点保存
  2. 声明式图定义:通过节点(Node)和边(Edge)抽象复杂的多步骤任务流
  3. 持久化支持:内置 PostgreSQL、MySQL、Oracle、MongoDB、Redis、文件系统等存储后端
  4. 可观测性集成:与 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.javaKeyStrategy.javastate/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;
    }
}

把策略下放到每个键上,好处归结为三点:

  1. 避免覆盖冲突:尤其在并行节点中,多个分支同时写状态不会互相覆盖
  2. 支持累积模式:对话历史、审批链这类需要不断追加的场景能自然表达
  3. 状态可预测:合并规则是声明式注册的,不依赖执行顺序,跑出来的结果稳定

使用示例:

KeyStrategyFactory factory = KeyStrategyFactoryBuilder.builder()
    .addDefault("chatHistory", NodeAggregationStrategy.APPEND)  // 聊天历史追加
    .addDefault("currentUser", NodeAggregationStrategy.REPLACE) // 用户信息替换
    .build();

核心设计 2:检查点与中断恢复系统

位置checkpoint/ 包、InterruptableAction.java

这是我最关心的部分,也是手写最难做对的地方。graph-core 支持可中断的工作流执行,关键组件有三个:

  1. CheckpointSaver:抽象持久化接口

    • 实现:PostgresSaverMysqlSaverRedisSaverMongoSaverFileSystemSaver
    • 保存当前节点状态、下一节点 ID、完整状态快照
  2. 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();
        }
    }
    
  3. 恢复机制

    • 保存检查点时分配唯一 ID(UUID)
    • 恢复时加载检查点,重建状态和节点位置
    • 支持时间旅行:回溯到任意历史检查点

检查点数据结构:

public class Checkpoint {
    private final String id;                    // 检查点 ID
    private Map<String, Object> state;          // 序列化的状态
    private String nodeId;                      // 当前节点
    private String nextNodeId;                  // 下一个节点
}

这套设计的价值,落到实际场景里就是三类需求:

  1. Human-in-the-Loop:在关键决策点暂停,等待人工审批
  2. 长时运行任务:支持任务跨会话持久化,服务重启不丢进度
  3. 调试与审计:可回溯任意执行点,定位问题不用靠猜

典型场景就是审批节点:

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.javaexecutor/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();
        });
    }
}

响应式带来的好处,归结为三点:

  1. 实时反馈:前端可立即显示中间节点结果,不用等整条流程跑完
  2. 资源效率:避免一次性把所有结果加载到内存
  3. 灵活组合:可与 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 的理由其实就几条:

  1. 已经是 Spring Boot 项目,技术栈一致
  2. 需要生产级的可观测性和监控
  3. 重视响应式流处理和实时反馈
  4. 需要企业级数据库持久化(PostgreSQL/Oracle)

总结

读完 graph-core,最深的体会是:一个工作流引擎的复杂度,不在流程控制本身,而在围绕流程的那套运行时能力。流程控制用循环也能写,但状态合并、检查点中断、响应式流式输出,这三块任何一个自己做都容易出 bug。

回头看,graph-core 三个核心设计对应的就是三个最难的运行时问题:

  1. KeyStrategy:状态怎么在节点(尤其并行节点)间安全传递
  2. CheckpointSaver + InterruptableAction:跑到一半怎么存档、怎么恢复、怎么等人审
  3. Reactor 执行引擎:结果怎么边算边吐、怎么背压、怎么并行汇聚

把这三块看明白,自己再写多步骤编排时,无论是直接用它,还是借鉴思路自己实现,心里都有底了。