BLOGE
0.9.8-RC1在线版 · 事实校验 2026-09-15 · English
第 22 章 —— 分布式运行与容量验证
本章承诺: 你会把同一条订单 graph 放到两个共享 durable store 的 worker 上,解释租约为何需要 fencing、路由为何必须稳定,并产出一份只对明确负载和环境成立的容量判断。
学习目标
学完本章,你能够:
- 解释 lease owner、timeout 和 epoch 如何完成一次安全接管。
- 区分 execution 路由稳定性与当前 shard resolver 的新选择。
- 找出扩容前真正饱和的资源,而不是默认增加 worker。
- 设计 arrival rate、service time、并发和 backlog 的最小容量实验。
- 写清一次测试证明的分布式合同及其外部边界。
这一章只增加一个变量
上一章把世界固定在一个 JVM。现在订单仍是同一条 graph,但 Worker A 与 Worker B 共享 durable store。新增的问题只有一个:执行的所有权不再由进程内存天然保证。
这会带来两种危险:A 已失联却再次写入;同一业务键因新路由决定漂到另一个 shard。
开场:A 超时后,B 能直接继续吗
先预测这段时间线是否安全:
10:00:00 Worker A claim order-42, epoch=1
10:00:20 A 写入 checkpoint
10:01:01 lease 过期;Worker B claim, epoch=2
10:01:03 A 网络恢复,再写一次
如果 store 只检查“owner 名字”,A 可能覆盖 B 的新状态。安全接管不仅需要 timeout,还需要单调递增的 fencing epoch。
运行两项 RC1 分布式合同测试
cd submodule/bloge
mvn -pl bloge-durable \
-Dtest=ClusterFailoverSimulationTest,RuntimeRoutingShardingTest test
2026-09-14 在 commit cc38fbe5 上观察到:
RuntimeRoutingShardingTest Tests run: 2, Failures: 0
ClusterFailoverSimulationTest Tests run: 1, Failures: 0
Tests run: 3, Failures: 0, Errors: 0, Skipped: 0
测试使用逻辑时钟和内存 store,不依赖真实等待。它建立的是 deterministic contract,不是数据库压测。
机制一:接管是一条三条件链
ClusterFailoverSimulationTest 做了四件事:
- A 认领 execution,得到 epoch 1,并在该上下文写 checkpoint。
- 逻辑时钟越过一分钟 timeout。
- B 的 recovery loop 重新认领,epoch 变为 2,写恢复标记并完成 execution。
- A 携带旧 epoch 1 再写,store 抛出
StaleFencingEpochException。
安全接管需要三个条件同时成立:
lease 已过期
+ 新 owner 原子认领并提升 epoch
+ 每次受保护写入校验当前 epoch
只有前两个条件时,僵尸 worker 仍能污染状态;只有第三个条件而没有可接管扫描,执行会永远停住。
单因素破坏:取消 stale epoch 校验
不改变 timeout、checkpoint 或恢复逻辑,只让 store 接受 epoch 1 的最后一次写入。最终状态可能仍显示 COMPLETED,但恢复标记已被旧 owner 覆盖。这说明“最终状态一样”不足以证明接管正确。
恢复动作是让所有改变受保护 execution 状态的 store 写入携带 lease context,并由持久化层 fail closed;不能靠应用日志事后猜冲突。
机制二:路由是业务键的稳定绑定
RuntimeRoutingShardingTest 首次执行时把:
tenant-a + order-999 → shard-a
写入 RoutingStore。第二个 engine 的 resolver 即使返回 shard-b,已有绑定仍胜出,第二次执行继续使用 shard-a。
resolver 是“没有绑定时去哪”的策略;routing store 是“这个业务对象已经在哪”的事实。若每次都重新 hash,扩容改变 shard 数量后,恢复和事件关联可能找不到旧状态。
迁移必须是一项显式控制面动作:冻结业务键、搬运状态、更新 binding、验证读写,再恢复流量。它不是普通 resolver 重新计算。
容量判断:先找瓶颈,再决定扩什么
把吞吐想成一条窄管道:
入口 arrival rate
→ worker 可运行并发
→ external API 配额
→ connection pool
→ durable store 写入与锁竞争
增加 worker 只扩大第二段。若 database checkpoint latency 已饱和,更多 worker 会增加排 队和锁竞争。
用 Little's Law 建立第一个估算:稳定系统中 in-flight ≈ arrival rate × average service time。它是测量设计的起点,不是 BLOGE 的容量保证。
| 要观察的量 | 说明 |
|---|---|
| arrival rate | 每秒新 execution / work item |
| service time | 从可运行到完成的分布,不只平均值 |
| in-flight | 当前占用 worker 或外部资源的数量 |
| backlog age | 最老待处理项等待多久 |
| store latency | checkpoint、claim、poll 的 p50/p95/p99 |
一个可归因的扩容实验
固定 graph、输入分布、数据库、连接池和外部 stub,只把 worker 从 1 增加到 2、4、8。每一步记录吞吐、p95、backlog age 和 store latency。
- 若吞吐近似增加且 store latency 稳定,worker 可能是当前瓶颈。
- 若吞吐不变而 store latency 上升,store/锁竞争更可疑。
- 若 backlog 下降但外部 429 增加,你只是把压力推给下游。
一次只改变一个轴,否则不能把变化归因给 worker 数量。
远端 worker 是派发边界,不是免费容量
WorkerDispatcher 把特定 work item 发到远端执行池;Kafka 模块提供一种 transport。派发后新增序列化、投递、重复消息、heartbeat、结果关联和取消语义。远端 GPU worker 可以解决硬件亲和性,却不自动解决 exactly-once effect。
正文只保留这个边界。topic、heartbeat、batch size、schema migration 状态和 benchmark 参数集中在运行参考,避免把架构判断埋在配置目录里。
现实映射:连锁仓库的订单接管
仓库 A 扫描了包裹却断网,仓库 B 不能只因为“一分钟没消息”就同时发货。总部要发出新一代作业凭证,并拒绝 A 的旧凭证。订单原先绑定华东仓,也不能因为今天新增华南仓就自动漂移。
这分别对应 fencing epoch 和 routing binding。增加仓库数量只有在拣货工位是瓶颈时才提升吞吐;如果瓶颈是中央库存锁,仓库越多争用越严重。