BLOGE
0.9.8-RC1在线版 · 事实校验 2026-09-15 · English
第 13 章 —— 持久执行
承诺: 读完本章后,你将理解 BLOGE 如何持久化执行状态,使 graph 在进程重启后恢复已接受的 checkpoint;当 graph definition 与 durable store 都可重建时,满足条件的另一 worker 可以继续执行。
学习目标
- 解释"持久执行"在 BLOGE 中的含义 —— 基于检查点的持久化,涵盖节点输出、挂起上下文、循环快照和执行标识,使 graph 能在崩溃后恢复。
- 识别使持久化 成为可能的运行时存储:
ExecutionStore、ExecutionCheckpointStore、WaitStore和WorkItemStore。 - 将内存持久化存储接入
GraphEngine,并在两个引擎实例之间运行挂起 → 信号 → 恢复流程。 - 描述冷恢复的工作方式:执行标识 → graph 注册表查找 → 检查点重新加载 → 继续执行。
- 识别检查点类型(
NODE_OUTPUT、LOOP_SNAPSHOT、FOREACH_PROGRESS、SUSPEND_CONTEXT、PARTIAL_RESULT)及各自的写入时机。
前置条件
- 第 6 章 —— 韧性设计 —— 重试和超时基础。
- 第 12 章 —— 等待世界的回应 —— 挂起、信号和事件关联。
- 熟悉
SuspendableOperator和OperatorResult.suspend()。
源示例
| 文件 | 展示内容 |
|---|---|
RuntimeGraphEngineDurabilityTest.java | 使用运行时存储的跨进程信号;两个引擎实例共享状态 |
RuntimeGraphRegistryRecoveryTest.java | 从已发布的 graph 定义进行冷恢复;hash 不匹配拒绝;过期租约自动恢复 |
PaymentWaitExample.java | 使用运行时存储的端到端挂起 → publishEvent → 恢复生命周期 |
payment-wait.bloge | DSL await 块,包含事件关联、超时和分支 |
RuntimeLoopOperatorDurabilityTest.java | 循环快照检查点和从迭代中恢复 |
为什么这很重要
到目前为止你学到的一切 —— 节点、依赖、分支、韧性、等待 —— 都假设 JVM 保持存活。在生产环境中,这个假设会失效:
- 进程在执行过程中崩溃。 没有持久化,所有进行中的节点输出都会丢失, 工作流必须从头重新开始。
- 挂起可能持续数分钟、数小时甚至数天(等待支付回调、经理审批)。让一个线程阻塞这么长时间会浪费资源,且任何部署都会将其销毁。
- 水平扩展 意味着多个 worker 进程。在 Worker A 上挂起的 graph 必须能在 Worker B 上恢复,而 Worker A 可能永远不会回来。
持久执行解决了这三个问题。引擎在节点完成时写入 检查点 到持久化存储,并在节点挂起时写入 等待记录。恢复时 —— 无论是由信号、定时器还是崩溃恢复触发 —— 任何引擎实例都可以重新加载这些检查点,并从执行暂停的精确前沿继续。
核心洞察: BLOGE 中的持久化不是“事件溯源”或“从头重放”。引擎会 跳过 checkpoint 被接受的节点,只调度剩余前沿。若外部副作用已发生、但对应 checkpoint 尚未提交,该结果仍可能处于未知状态;幂等与对账仍是应用责任。
心智模型
把持久执行想象成 食谱书中的书签:
当支付事件到达时,任何引擎实例都可以:
- 从
ExecutionStore加载ExecutionInstance。 - 重新加载 三个已完成节点的
NODE_OUTPUT检查点。 - 解析
WAIT_SIGNAL等待记录。 - 从
awaitPayment恢复 → branch →fulfillOrder。
该模型的四大支柱:
| 支柱 | 存储 | 记录内容 |
|---|---|---|
| 标识 | ExecutionStore | 执行生命周期、状态、租约、恢复尝试次数 |
| 检查点 | ExecutionCheckpointStore | 节点输出、循环快照、挂起上下文、foreach 进度 |
| 等待 | WaitStore | 执行为何暂停(信号、定时器、任务、事件) |
| 工作项 | WorkItemStore | 待处理操作:定时器到期、事件匹配、任务恢复 |
执行生命周期如果画成状态图,会长这样:
bloge-runtime-spi 模块
在写任何代码之前,先给"缝"命名。持久化执行层被切成两个模块:
bloge-runtime-spi—— 只放接口。ExecutionStore、ExecutionCheckpointStore、WaitStore、TimerService、GraphRegistryStore、LeaseStore、AuditJournalStore、TaskInboxStore,加上这些 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,并且已经完成的
工作仍保持完成。
跟随 loan-42,不要跟随 engine 对象
| 时刻 | durable 事实 | 运行后果 |
|---|---|---|
| Engine 1 写 checkpoint | fetchApplication 与 checkCredit 已完成 | 两者输出可以恢复 |
| 进程停止 | heap 对象全部消失 | 内存引用不能证明任何 durable 事实 |
Engine 2 claim loan-42 | graph 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_CONTEXT | operator 返回 OperatorResult.suspend() 时 | 挂起键、部分输出、graph 上下文快照 |
PARTIAL_RESULT | 与 SUSPEND_CONTEXT 一起写入 | 挂起期间可见的中间输出 |
LOOP_SNAPSHOT | 每次循环迭代后 | completedIterations + carryStateJson,以便循环从中间恢复 |
FOREACH_PROGRESS | 每个顺序 foreach 项之后 | 逐项进度,以便 foreach 可以跳过已完成的项 |
EVENT_CORRELATION | await 注册匹配器时 | 关联键,以便事件路由到正确的执行 |
SESSION_PHASE_SNAPSHOT / SESSION_ROUND_SNAPSHOT | Session 的 phase / 轮次转换时 | Session 特有的状态(参见第 14 章 —— 多轮 Session) |
执行生命周期状态
ExecutionStatus 跟踪执行在其生命周期中的位置:
RUNNING→ 节点正在活跃执行SUSPENDED→ 至少一个节点返回了OperatorResult.suspend()COMPLETED→ 所有终端节点成功完成FAILED→ 节点在所有韧性尝试后失败FAILED_RECOVERY→ 自动崩溃恢复超过了maxRecoveryAttemptsCANCELLED/TERMINATED→ 被显式停止