Skip to main content

BLOGE 0.9.8-RC1 在线版 · 事实校验 2026-09-15 · English

第 13 章 —— 持久执行

承诺: 读完本章后,你将理解 BLOGE 如何持久化执行状态,使 graph 在进程重启后恢复已接受的 checkpoint;当 graph definition 与 durable store 都可重建时,满足条件的另一 worker 可以继续执行。


学习目标

  1. 解释"持久执行"在 BLOGE 中的含义 —— 基于检查点的持久化,涵盖节点输出、挂起上下文、循环快照和执行标识,使 graph 能在崩溃后恢复。
  2. 识别使持久化成为可能的运行时存储:ExecutionStoreExecutionCheckpointStoreWaitStoreWorkItemStore
  3. 将内存持久化存储接入 GraphEngine,并在两个引擎实例之间运行挂起 → 信号 → 恢复流程。
  4. 描述冷恢复的工作方式:执行标识 → graph 注册表查找 → 检查点重新加载 → 继续执行。
  5. 识别检查点类型(NODE_OUTPUTLOOP_SNAPSHOTFOREACH_PROGRESSSUSPEND_CONTEXTPARTIAL_RESULT)及各自的写入时机。

前置条件

源示例

文件展示内容
RuntimeGraphEngineDurabilityTest.java使用运行时存储的跨进程信号;两个引擎实例共享状态
RuntimeGraphRegistryRecoveryTest.java从已发布的 graph 定义进行冷恢复;hash 不匹配拒绝;过期租约自动恢复
PaymentWaitExample.java使用运行时存储的端到端挂起 → publishEvent → 恢复生命周期
payment-wait.blogeDSL await 块,包含事件关联、超时和分支
RuntimeLoopOperatorDurabilityTest.java循环快照检查点和从迭代中恢复

为什么这很重要

到目前为止你学到的一切 —— 节点、依赖、分支、韧性、等待 —— 都假设 JVM 保持存活。在生产环境中,这个假设会失效:

  • 进程在执行过程中崩溃。 没有持久化,所有进行中的节点输出都会丢失,工作流必须从头重新开始。
  • 挂起可能持续数分钟、数小时甚至数天(等待支付回调、经理审批)。让一个线程阻塞这么长时间会浪费资源,且任何部署都会将其销毁。
  • 水平扩展 意味着多个 worker 进程。在 Worker A 上挂起的 graph 必须能在 Worker B 上恢复,而 Worker A 可能永远不会回来。

持久执行解决了这三个问题。引擎在节点完成时写入 检查点 到持久化存储,并在节点挂起时写入 等待记录。恢复时 —— 无论是由信号、定时器还是崩溃恢复触发 —— 任何引擎实例都可以重新加载这些检查点,并从执行暂停的精确前沿继续。

核心洞察: BLOGE 中的持久化不是“事件溯源”或“从头重放”。引擎会 跳过 checkpoint 被接受的节点,只调度剩余前沿。若外部副作用已发生、但对应 checkpoint 尚未提交,该结果仍可能处于未知状态;幂等与对账仍是应用责任。


心智模型

把持久执行想象成 食谱书中的书签

Diagram: 13-durable-execution figure 1

当支付事件到达时,任何引擎实例都可以:

  1. ExecutionStore 加载 ExecutionInstance
  2. 重新加载 三个已完成节点的 NODE_OUTPUT 检查点。
  3. 解析 WAIT_SIGNAL 等待记录。
  4. awaitPayment 恢复 → branch → fulfillOrder

该模型的四大支柱:

支柱存储记录内容
标识ExecutionStore执行生命周期、状态、租约、恢复尝试次数
检查点ExecutionCheckpointStore节点输出、循环快照、挂起上下文、foreach 进度
等待WaitStore执行为何暂停(信号、定时器、任务、事件)
工作项WorkItemStore待处理操作:定时器到期、事件匹配、任务恢复

执行生命周期如果画成状态图,会长这样:

Diagram: 13-durable-execution figure 2


bloge-runtime-spi 模块

在写任何代码之前,先给"缝"命名。持久化执行层被切成两个模块:

  • bloge-runtime-spi —— 只放接口ExecutionStoreExecutionCheckpointStoreWaitStoreTimerServiceGraphRegistryStoreLeaseStoreAuditJournalStoreTaskInboxStore,加上这些 store 发出的强类型事件。
  • bloge-durable —— 引擎管道,消费这些接口:恢复循环、哈希校验、 checkpoint 编解码,以及测试用的内存实现。

数据库实现各自有模块(MyBatis/JDBC 对应 bloge-durable-mybatis)。这个 切分意味着:如果你只是想插入一个自定义的存储,依赖 bloge-runtime-spi 即可,无需引入引擎工件。

<dependency>
<groupId>com.leanowtech</groupId>
<artifactId>bloge-runtime-spi</artifactId>
</dependency>

第一个可运行的示例

本示例摘自 RuntimeGraphEngineDurabilityTest。它演示了一个 graph 在一个引擎实例上挂起、在另一个实例上完成 —— 这正是你在生产中需要的"跨进程信号"模式。

步骤 1 — 创建存储和 graph

// Shared durable stores (in production these would be database-backed)
var executionStore = new InMemoryExecutionStore();
var checkpointStore = new InMemoryExecutionCheckpointStore();
var waitStore = new InMemoryWaitStore();
var timerService = new InMemoryTimerService();

// A suspendable operator that parks execution
SuspendableOperator<Void, String> waitOp = (in, ctx) ->
OperatorResult.suspend("ck", null, Duration.ofMinutes(10));
Operator<Object, String> completeOp = (in, ctx) -> "C";

var graphBuilder = new GraphBuilder("runtime-cross-process");
var graph = graphBuilder
.suspendNode("w", waitOp)
.node("c", completeOp).dependsOn("w")
.build();

CountDownLatch suspended = new CountDownLatch(1);
AtomicReference<String> executionId = new AtomicReference<>();

ExecutionListener listener = new ExecutionListener() {
@Override
public void onGraphStart(String graphName, GraphContext ctx) {
executionId.set((String) ctx.get(ReservedKeys.EXECUTION_ID));
}

@Override
public void onNodeSuspended(String graphName, String nodeId, String suspendKey) {
if ("w".equals(nodeId)) {
suspended.countDown();
}
}
};

步骤 2 — 在引擎 1 上执行直到挂起

var engine1 = GraphEngine.builder()
.registry(new DefaultOperatorRegistry())
.executionStore(executionStore)
.executionCheckpointStore(checkpointStore)
.waitStore(waitStore)
.timerService(timerService)
.listeners(List.of(listener))
.build();

// 源测试把 executeWithOperators() 放到后台线程中运行,这样它就能先观察挂起,
// 再由第二个引擎实例恢复。
Thread runThread = Thread.ofVirtual().start(() -> {
try {
engine1.executeWithOperators(graph, null, graphBuilder.operators());
} catch (Exception ignored) {
}
});

// The execution is now SUSPENDED in the store
assertTrue(suspended.await(5, TimeUnit.SECONDS));
assertNotNull(executionId.get());
assertEquals(ExecutionStatus.SUSPENDED,
executionStore.get(executionId.get()).orElseThrow().status());

步骤 3 — 在不同的引擎 2 上恢复

var engine2 = GraphEngine.builder()
.registry(new DefaultOperatorRegistry())
.executionStore(executionStore) // same stores
.executionCheckpointStore(checkpointStore)
.waitStore(waitStore)
.timerService(timerService)
.build();

// Signal resumes the suspended node and continues to "c"
engine2.signal(graph, executionId.get(), "w", "resume-data");

// Execution completes
assertEquals(ExecutionStatus.COMPLETED,
executionStore.get(executionId.get()).orElseThrow().status());
assertEquals("C",
CheckpointCodec.DEFAULT.deserialize(
checkpointStore.load(executionId.get(), CheckpointType.NODE_OUTPUT, "c")
.orElseThrow().payload()));
runThread.join(1_000);

引擎 1 和引擎 2 是完全独立的对象。除了持久化存储外,它们不共享任何东西。这就是基本的生产模式:一个进程挂起,另一个进程恢复。


冷恢复:进程消失后,什么还能留下

warm resume 很容易被高估,因为原 engine 仍持有内存对象。真正的 durable 声明需要 Engine 1 消失、Engine 2 读取同一 executionId,并且已经完成的 工作仍保持完成。

图:冷恢复与未知 effect 边界

跟随 loan-42,不要跟随 engine 对象

时刻durable 事实运行后果
Engine 1 写 checkpointfetchApplicationcheckCredit 已完成两者输出可以恢复
进程停止heap 对象全部消失内存引用不能证明任何 durable 事实
Engine 2 claim loan-42graph hash 与 checkpoint 匹配只有未完成 node 可进入调度
恢复完成新 checkpoint 记录后续工作同一逻辑 execution 继续前进

identity 属于 execution record,不属于任一 engine instance。派发未完成工作 之前,系统会重新加载 graph definition,并与保留状态核对。

effect 仍可能是 UNKNOWN

设想 Engine 1 调用了 approveLoan,放款方已经提交批准,但进程在 completion checkpoint 之前停止。Engine 2 只看见一个未完成 node;checkpoint 无法回答 外部 effect 是否发生。

因此,恢复仍需要第 6 章的合同:稳定幂等键、reconciliation 查询,以及明确的 retry 或人工修复策略。durable execution 能在完成记录已保留时避免再次调度 已完成的引擎工作;它不能为一个未协调的外部系统凭空制造 exactly-once effect。


拆解分析

检查点类型

引擎在不同的生命周期时刻写入不同的检查点类型。这些定义在 CheckpointType 中:

类型写入时机用途
NODE_OUTPUT每个节点完成时序列化的输出 —— 恢复时跳过
SUSPEND_CONTEXToperator 返回 OperatorResult.suspend()挂起键、部分输出、graph 上下文快照
PARTIAL_RESULTSUSPEND_CONTEXT 一起写入挂起期间可见的中间输出
LOOP_SNAPSHOT每次循环迭代后completedIterations + carryStateJson,以便循环从中间恢复
FOREACH_PROGRESS每个顺序 foreach 项之后逐项进度,以便 foreach 可以跳过已完成的项
EVENT_CORRELATIONawait 注册匹配器时关联键,以便事件路由到正确的执行
SESSION_PHASE_SNAPSHOT / SESSION_ROUND_SNAPSHOTSession 的 phase / 轮次转换时Session 特有的状态(参见第 14 章 —— 多轮 Session

执行生命周期状态

ExecutionStatus 跟踪执行在其生命周期中的位置:

  • RUNNING → 节点正在活跃执行
  • SUSPENDED → 至少一个节点返回了 OperatorResult.suspend()
  • COMPLETED → 所有终端节点成功完成
  • FAILED → 节点在所有韧性尝试后失败
  • FAILED_RECOVERY → 自动崩溃恢复超过了 maxRecoveryAttempts
  • CANCELLED / TERMINATED → 被显式停止

等待类型

WaitType 记录执行 为何 被挂起:

  • WAIT_SIGNAL — 等待具有匹配键的 engine.signal()
  • WAIT_TIMER — 基于超时的自动恢复
  • WAIT_EVENT — 等待具有关联匹配的 engine.publishEvent()
  • WAIT_TASK — 等待用户任务完成
  • WAIT_PHASE_TIMEOUT — session 的 phase 级超时已触发
  • WAIT_ROUND_TIMEOUT — session 的 round 级超时已触发
  • WAIT_RETRY_BACKOFF — 作为持久工作项调度的延迟重试

如果 session 对你还是新主题,可以阅读第 14 章 —— 多轮 Session, 那里会把 phase / round 超时和 session 快照如何建立在同一套 durable 机制之上讲清楚。状态机的持久化则走了另一条路径:StateMachineCheckpoint 通过通用 execution-checkpoint SPI 接入,而不是再定义一种新的核心检查点类型;你会在第 15 章 —— 状态机里看到它。

DurableManager 桥接层

DurableManagerGraphEngine 用于协调跨所有存储写入的内部桥接层。你不需要直接调用它 —— 引擎在以下时机调用它:

  • Graph 启动 → initExecution()
  • 节点完成 → saveNodeCompletion()
  • 节点挂起 → 写入挂起上下文以及对应的等待/定时器记录
  • 循环迭代完成 → saveLoopSnapshot()
  • 恢复完成 → 清理挂起上下文 + 等待记录

常见陷阱

❌ 假设检查点数据保留其原始 Java 类型

当你挂起并稍后恢复时,检查点载荷会被序列化为 JSON 然后反序列化回来。如果你的 operator 输出是一个 PaymentEvent 记录,反序列化后它将变成 Map<String, Object>,除非你使用了类型化的 CheckpointCodec

// FRAGILE — will ClassCastException after resume
PaymentEvent event = results.get("awaitPayment", PaymentEvent.class);

// SAFE — handle both fresh-run and checkpoint-restored shapes
Object raw = results.getRaw("awaitPayment");
String txnId;
if (raw instanceof PaymentEvent pe) {
txnId = pe.transactionId();
} else if (raw instanceof Map<?,?> m) {
txnId = (String) m.get("transactionId");
}

这种模式在 PaymentWaitExample 中随处可见。

修复方法: 要么使用 getRaw() 配合模式匹配,要么将 JacksonCheckpointCodec 的类型化反序列化接入 GraphEngine.builder().checkpointCodec(...)


常见故障

schema 变更后,检查点反序列化失败

你发布了 operator 的 v2 版本,把输出类型从 record OrderOutput(String orderId, int total) 改成了 record OrderOutput(String orderId, BigDecimal total, String currency)

某个还在运行中的执行使用的是 v1 schema 挂起并写入了检查点。恢复时加载这个检查点,就会遇到这样的失败:

Recovery error: cannot deserialize the checkpoint for node 'createOrder'
because required field 'currency' is missing from the stored payload.

修复: schema 变更尽量保持向后兼容(新增可选字段,不要删除必填字段)。如果确实是 breaking change,就先清空在途执行,或者实现一个能处理版本迁移的 CheckpointCodec


引导式重写

打开 RuntimeGraphRegistryRecoveryTest。测试 publishesDslGraphAndRecoversColdSignalWithoutCallerGraph 演示了 冷恢复 —— 在调用方不传入 Graph 对象的情况下恢复:

// Engine 1: execute with a DSL graph — the engine publishes the
// GraphDefinition to graphRegistryStore automatically
GraphEngine engine1 = GraphEngine.builder()
.registry(registry)
.executionStore(executionStore)
.executionCheckpointStore(checkpointStore)
.waitStore(waitStore)
.graphRegistryStore(graphRegistryStore) // ← enables cold recovery
.addGraphDefinitionCodec(codec) // ← how to decode DSL
.timerService(timerService)
.build();

engine1.execute(graph, new GraphContext(Map.of("sessionId", "SESSION-001")));
// Graph suspends at "wait" node

// Engine 2: signal WITHOUT passing the graph object
GraphEngine engine2 = GraphEngine.builder()
.registry(recoveryRegistry)
.executionStore(executionStore)
.executionCheckpointStore(checkpointStore)
.waitStore(waitStore)
.graphRegistryStore(graphRegistryStore)
.addGraphDefinitionCodec(codec)
.timerService(timerService)
.build();

// The engine loads the graph from graphRegistryStore using the
// execution's graphVersion + graphHash binding
engine2.signal(executionId, "wait", Map.of("message", "resume"));

思考问题:

  1. 如果 graph 定义在挂起和恢复之间发生了变化会怎样? 引擎会比较 ExecutionIdentity 中存储的 graphHash 与已发布 GraphDefinition 的 hash。如果不匹配,恢复会以 "hash mismatch" 错误失败。这在 rejectsColdRecoveryWhenPublishedGraphHashChanges 测试用例中有测试。

  2. 自动崩溃恢复是怎样的? 当启用了 RecoveryConfig 并且某次执行的租约过期(worker 崩溃了)时,另一个引擎实例可以自动认领过期的执行并重新运行它。recoveryAttempts 计数器每次递增。如果预算耗尽,状态转变为 FAILED_RECOVERY

  3. Graph 定义存储在哪里?GraphRegistryStore 中,它持久化 graphName + graphVersion + graphHash + source。在生产环境中,这由 bd_graph_definition 数据库表支撑。


Graph 定义来源(GraphDefinitionSource)

当引擎把 graph 发布到 GraphRegistryStore 时,它需要知道这个定义是从哪里来的 —— 是磁盘上的 DSL 文件、Java builder 调用,还是从 registry 导入的。GraphDefinitionSource 接口(com.leanowtech.bloge.core.runtime.registry.GraphDefinitionSource)正是用来捕获这个来源信息的:

public interface GraphDefinitionSource {
String sourceId(); // 该来源的唯一标识
String sourceType(); // 例如 "dsl-file"、"java-builder"、"registry"
Optional<String> content(); // 如果可用,保存原始源内容
}

在构建 graph 时附加来源信息,用于 registry 发布:

var source = new FileBasedGraphDefinitionSource(
"file:///graphs/order-process.bloge", dslContent);

Graph graph = Graph.builder("orderProcess")
.definitionSource(source)
// ... nodes and edges ...
.build();

engine.publishGraphDefinition(graph);

三个字段服务于不同的恢复需求:

字段用途
sourceId()唯一标识来源 —— 文件路径、classpath 资源或 registry URI
sourceType()告诉 registry 定义是如何产生的("dsl-file""java-builder""registry"
content()可选地存储原始 DSL 文本,这样即使没有原始文件也能重建定义

如果没有 GraphDefinitionSource,registry 只会存储编译后的 graph 结构,丢失回到原始构件的线索。对于需要审计追踪或多集群复制的生产系统,务必附加来源信息。


检查点编解码器 SPI(CheckpointCodec)

默认情况下,引擎将检查点载荷序列化为 JSON。对于很多工作流这已经够用了,但生产系统有时需要:

  • 自定义压缩 —— 降低高吞吐工作流的存储成本
  • 加密 —— 保护包含敏感数据(PII、财务记录)的检查点
  • 兼容性层 —— 处理跨部署的 schema 迁移

CheckpointCodec SPI(com.leanowtech.bloge.core.checkpoint.CheckpointCodec)让你可以替换默认序列化方式:

public interface CheckpointCodec {
byte[] encode(ExecutionCheckpoint checkpoint);
ExecutionCheckpoint decode(byte[] data);
}

在引擎 builder 中接入:

GraphEngine engine = GraphEngine.builder()
.registry(registry)
.checkpointCodec(new JsonCheckpointCodec()) // 自定义序列化
.build();

类型化的 codec 还能消除上面"常见陷阱"部分描述的 Map<String, Object> 问题 —— 如果 codec 知道目标类型,恢复后的检查点就能直接还原为正确的 Java record,而不是裸 map。


Operator 指纹不匹配策略

冷恢复仍会拒绝 graph definition hash 不匹配。另一个独立问题是:完成节点保存了 Operator 指纹,但当前注册的 Operator 已换成另一指纹。VersionMismatchPolicy 只控制这类 checkpoint 决策:

策略行为
WARN恢复 checkpoint 并记录 Operator 不匹配;这是默认值。
RERUN跳过该 checkpoint,让节点重新进入调度。
FAILOperatorVersionMismatchException 中止恢复。

在引擎 builder 中明确配置:

GraphEngine engine = GraphEngine.builder()
.registry(registry)
.versionMismatchPolicy(VersionMismatchPolicy.FAIL)
.build();

选择依据应是副作用语义,而不是环境标签。只有节点可安全重复、不会复制不可逆副作用时,RERUN 才成立。WARN 会让新代码接受旧输出,因此需要显式兼容性判断。FAIL 会停止自动恢复,把迁移或人工处置留在引擎之外。

生产提醒——schema migration 是另一条版本轴

bloge-durable-mybatis artifact 自带、按方言区分的 Flyway migration 为事实源。/actuator/bloge/schema 会报告 UP_TO_DATEPENDINGFAILED;依赖新 schema 的 binary 上线前,应先应用所需 migration。RC1 的 V25 只迁移旧 session enum 值;recovery lease 来自 V17,audit journal 表来自 V9,lease fencing 来自 V26。运维迁移清单集中在附录 H,不埋在叙事章节里。


思维检查

  1. 使持久执行成为可能的四个运行时存储是什么? (ExecutionStore、ExecutionCheckpointStore、WaitStore、WorkItemStore。)

  2. 当引擎恢复一个挂起的执行时,已完成的节点会怎样? (它们被跳过。引擎加载它们的 NODE_OUTPUT 检查点并填充 NodeResults,而不重新执行 operator。)

  3. 为什么引擎在 graph hash 变化时拒绝冷恢复? (因为 graph 结构可能已经改变 —— 新节点、删除的边、不同的分支 —— 所以检查点的节点输出可能不再对应当前的 graph。恢复将产生未定义行为。)

  4. 当自动恢复超过 maxRecoveryAttempts 时,执行转变为什么 ExecutionStatusFAILED_RECOVERY —— 一个终态。)

  5. 如果一个 SuspendableOperator 返回 OperatorResult.suspend("key", partial, Duration.ofMinutes(5)),会创建多少条等待记录? (两条:一条 WAIT_SIGNAL 用于手动信号恢复路径,一条 WAIT_TIMER 用于 5 分钟超时自动恢复路径。)

  6. 设计题: 你的持久工作流会处理 2–4 周的保险理赔。在这段时间里,系统会多次发版。你准备如何确保这些在途执行能够穿越部署继续恢复?(保持 operator 输出 schema 的向后兼容;使用 GraphRegistryStore 给 graph 做版本化,让恢复流程能找到正确的 graph 定义;把 compensation operator 设计成幂等的,这样恢复后即使某个节点重跑也安全;并且在正式发布前,用长时间运行的执行在 staging 中做恢复演练。)


实验

目标: 构建一个两阶段持久支付流程,并模拟阶段之间的崩溃。

  1. 创建一个 graph,包含三个节点:

    • createOrder —— 一个返回订单 ID 的普通 operator。
    • awaitPayment —— 一个 SuspendableOperator,使用 OperatorResult.suspend("payment:" + orderId, null, Duration.ofSeconds(30)) 挂起。
    • fulfillOrder —— 依赖 awaitPayment,返回一个物流 ID。
  2. 接入运行时存储(使用内存实现):

    var executionStore = new InMemoryExecutionStore();
    var checkpointStore = new InMemoryExecutionCheckpointStore();
    var waitStore = new InMemoryWaitStore();
  3. 执行 graph 在 engine1 上。验证:

    • 执行状态为 SUSPENDED
    • createOrder 存在 NODE_OUTPUT 检查点。
    • awaitPayment 存在 SUSPEND_CONTEXT 检查点。
    • 等待记录包含 WAIT_SIGNALWAIT_TIMER
  4. 创建 engine2,使用相同的存储但全新的 DefaultOperatorRegistry。对挂起的节点发送信号:

    engine2.signal(graph, executionId, "awaitPayment", paymentData);
  5. 断言

    • 执行完成。
    • fulfillOrder 有一个 NODE_OUTPUT 检查点。
    • 所有等待记录已被清理。

进阶: 启用 GraphRegistryStore,尝试在不传入 graph 对象的情况下调用 engine2.signal() —— 仅使用执行 ID。


实验验收卡

  • 预期与观察: 新进程凭 identity、definition 和 checkpoint 接管。
  • 失败与恢复: 改变 definition version;恢复匹配版本或显式迁移。
  • 证明边界: 证明引擎冷恢复,不证明外部 effect 状态已知。
  • 练习合同: 中断 checkpoint;只改版本;交付接受和拒绝记录;完成节点不重跑且 UNKNOWN 单列即停止。

回顾

  • 持久执行 意味着持久化节点输出、挂起上下文、等待记录和执行标识,使 graph 能在崩溃后存活,并由任何引擎实例恢复。
  • 四个运行时存储 —— ExecutionStoreExecutionCheckpointStoreWaitStoreWorkItemStore —— 是基础。开发时使用内存实现;生产环境使用 MyBatis 支撑的实现。
  • 恢复时,引擎 加载检查点并跳过已完成的节点。它不会重放已完成的工作。
  • 冷恢复 使用 GraphRegistryStore 通过 graphVersion + graphHash 绑定重新加载 graph 定义 —— 不需要调用方提供 Graph 对象。
  • 引擎在 graph hash 变化 时拒绝恢复,防止检查点/graph 不匹配。
  • 自动崩溃恢复 使用租约过期和 RecoveryConfig,让健康的 worker 认领过期的执行,并带有 recoveryAttempts 预算。
  • 检查点载荷是 JSON 序列化的。使用 getRaw() 配合模式匹配,或自定义 CheckpointCodec 来处理全新运行与恢复执行之间的类型差异。

下一步

第 14 章 —— 多轮 Session 中,你将从“持久的一次性 graph”走向“更长生命周期的交互”。你在这里学到的检查点与恢复机制,会直接成为 phase / round session 的基础,并继续延伸到第 15 章 —— 状态机


参考链接

Coding Agent: Open the versioned task guide.