系统设计 / 任务运行时 · 从一次 run() 到可观测的调度层 / 同 key 协调:重复调用的去向 待审核 5 / 10
coordinator

同 key 协调:重复调用的去向

搜索框每敲一个字发一次请求,保存按钮被连点三下,同一份数据在两个组件里各拉一次。这些调用的共同点是带着同一个身份:同一个 URL、同一个文档 id、同一句查询。运行时默认把它们当作互不相干的工作,各自入队各自执行。真实的期望却因场景而异:搜索要最后一次,保存要按顺序一笔一笔写,取数据只想发一次请求。@vega/job 把这个决定从内核里挪了出去,交给一个叫 coordinator 的小函数。

1 · 同一个 key 上的重复调用

run(runtime, taskFn) 不认识「重复」。同一个 taskFn 连提交 4 次就是 4 个任务,并发额度够用时它们同时开跑,taskFn 体执行 4 次,四次调用各拿各的结果(core/scenarios.test.ts §1)。

要改变这一点,先得给调用一个身份。coordinator 是一个函数,运行时在派发每个带 key 的调用之前先把它交给这个函数,由它决定新调用是照常执行、并进某个在飞的同 key 任务,还是干脆不执行。withCoordinator(taskFn, coordinator, { key }) 把策略与身份一起绑在 taskFn 上,run(runtime) 收到这样的 taskFn 就走协调路径。

协调器眼里有两个角色:本次调用的任务叫 intent,同 key 的其他在飞任务叫 peer。内建策略之间的差别全部落在「peer 怎么办」这一句上。

警示 · key 不在 RunOptions 里,协调也不在内核里。key 走 withCoordinator 的第三个参数 withCoordinator(taskFn, coordinator, { key }),缺了它在绑定时就抛错,而不是等到 run() 才发现。运行时还必须写成 createRuntime({ plugins: [installCoordination] }):不装这个插件,被 withCoordinator 包过的 taskFn 一进 run() 就撞上内核的 loud-fail 守卫,整批任务 error。协调是可选装的一层,这两处就是它与内核之间仅有的接缝。

图 1-1 · 同一个 key 上连发 4 次调用,泳道每行一次调用。可切换五种协调器并调节 taskFn 体耗时,观察 taskFn 执行次数、各调用的终态与 cancelReason。

2 · dedupe 的合并

dedupeCoordinator 的实现只有两行:在 peer 里找一个非终态的,找到就 ctx.reuse(peer),找不到就 ctx.proceed()。同一个 key 上连发 4 次,taskFn 体只执行 1 次;后 3 次调用的任务以 cancelled 收尾、reason 是 dedupe,而它们的 promise 全部 resolve 到第一次调用的那个结果(core/scenarios.test.ts §2)。

「任务 cancelled 而 promise 成功」是这一档最容易读错的地方。task 记的是这次调用有没有自己执行,promise 记的是调用方要的值有没有拿到,两者本就不必一致。Go 的 singleflight 与 TanStack Query 的 request deduplication 是同一个语义。RxJS 的 share 不是——那是多播,一条流分给多个订阅者,不丢弃任何执行。

去重也不等于缓存:领头的任务一 settle,桶就空了,下一次同 key 调用重新起飞。要的若是「一段时间内不再请求」,那属于缓存层的职责,与协调器无关。

3 · restart 的接管

restartCoordinator 反过来做:先把所有非终态的 peer cancel() 掉,再让自己 proceed()。这是 RxJS switchMap 与 TanStack Query 的取消语义,search-as-you-type 的默认选择。

4 次连发的结局是前 3 次 cancelled、最后一次 success,而 taskFn 体仍然执行了 4 次(core/scenarios.test.ts §3)。执行次数不下降是它与 dedupe 的分水岭:每一次调用都真的起飞过,只是先到的没能跑完。取消是协作式的,被取消的 taskFn 要 await 到下一个 cancellation point 才真正退出,新调用不等它退出就入队(见任务与状态机 §3)。

有一处实测才看清的事:被 restart 取消的调用,cancelReasonabort,与用户手动 task.cancel() 得到的一模一样。「谁取消了我」这个问题在终态里读不出来,只能靠 coordination:restart 事件的 cancelledPeerIds。被取消的 taskFn 抛出的错误也不由运行时决定:本页 lab 里那个 sleep(ms, signal) 抛的是 AbortError,换成别的 cancellation point 就是别的错误对象。

4 · serialize 的队列

createSerializeCoordinator() 是唯一有状态的内建策略,每个 key 一条 FIFO 队列。首个调用直接 proceed,其后每个调用先 intent.parkCoordinator() 把自己停进 parked:coordinator,再排进队列,等队头 settle 后由队列放行。

4 次连发全部执行且全部成功,四个结果按调用次序依次产生(core/scenarios.test.ts §4)。它是唯一既不丢调用也不合并结果的策略,代价是总时长变成四段串行。适用场景是同一份文档的连续保存:每一笔都要写,且必须按提交次序落盘。

有状态就有生命期。工厂返回的 dispose() 得接到自己的销毁路径上,否则队列里的等待者永远等不到那个 reject。工厂另外返回的 resetAdmission() 则通常不用自己调——它会被自动挂进 backend 的 pre-cancel 钩子,cancelAll(runtime) 时顺带触发,把当前这代队列作废掉。

5 · exhaust 的丢弃

exhaustCoordinator 与前面三种一样从 @vega/job 直接导出,规则却是四条里最短的:只要有一个同 key 的 peer 非终态,新调用原地丢弃,peer 不受任何影响。它和 dedupe 一样无状态——ctx.peers() 已经按 hash 过滤过了,判断「有没有 peer 在飞」就是全部逻辑。RxJS 的 exhaustMap 是同一回事,redux-saga 的 throttle 去掉时间窗后也是。

同一个场景下它与 dedupe 的执行次数同为 1 次,结局却相反:后 3 次调用不但 cancelled(reason 为 exhaust-dropped),promise 还各自 reject 一个 CancelledError,拿不到 peer 的值(lab/facts.test.ts §5)。dedupe 把 4 次调用的命运绑在一起,exhaust 把它们彻底切开,peer 失败也牵连不到被丢弃的那几次。防双击的提交按钮要的正是后者:第二下点击应当什么都不发生,而不是跟着第一下一起报成功。

建议 · 四种策略对同一个问题给出四种答案,选哪一种取决于重复调用之间是否等价。等价且只读,用 dedupe;后来的取代先前的,用 restart;每一笔都要写且有序,用 serialize;后来的应当被忽略,用 exhaust。若连发本身就该被压掉,那是时间窗口 shaper的活,debounce 在窗口内合并调用,压根不会有第二个 intent 进到协调器。

6 · 协调协议

上面这些策略都由同一个协议写成。入口是 defineCoordinator(name, impl)name 不能为空,它会被自动盖在该协调器发出的每一条诊断事件上,devtools 靠它显示标签。

impl 收到的 CoordinationContext 就是协调器能碰的全部东西:

  • taskFn / options / hash:本次请求。hash 是 key 的散列,同 hash 才算同 key。
  • peers():同 hash 的其他在飞任务快照,不含 intent 自己。
  • intent:本次调用的任务句柄,进来时是 idle。
  • proceed():把 intent 推进内核的默认管线(等依赖、入队、执行),返回 taskFn 的结果。
  • reuse(other):把 intent 以 dedupe skip 掉,返回 peer 的 promise。
  • emit(event):往运行时主事件总线发一条 coordination:* 诊断。

proceed()reuse() 互斥,第二次调用直接抛错。协议要求每个 intent 恰好走其中一条;两条都没走时,内核的孤儿守卫兜底:发一条 coordination:orphaned-intent,再把 intent skip 成 cancelled,reason 为 orphaned-intent

诊断事件共四条:coordination:reusereuse() 发;coordination:restart 由 restart 自己发,带上被取消的 peer id;coordination:queued 由 serialize 在停靠等待者时发;coordination:orphaned-intent 由孤儿守卫发。它们与 task:* 生命周期事件走同一条总线,一次 subscribe(runtime) 全收(事件的两个类别见事件流与 devtools)。

图 6-1 · 两个手写 coordinator:bounded 把同 key 在飞数封顶,forgetful 什么都不做。可切换策略与 limit,观察日志里 coordination:* 事件的 coordinator 名字与各调用终态。

自定义策略的成本因此很低。把 dedupe 的「找一个 peer」改成「数一数 peer」,就得到一个把同 key 在飞数封顶的 bounded(limit):在飞数不足 limit 就 proceed,否则复用最早的那个 peer。limit 取 2 时 4 次连发执行 2 次,后两次以 dedupe 并进第一次的结果;limit 取 1 时它与内建 dedupe 的结局表逐格相同(lab/facts.test.ts §6)。

孤儿守卫的兜底同样可以跑出来核对。一个什么都不做的协调器让 4 次调用全部 cancelled、taskFn 体一次都没执行,而 4 个 promise 照常 resolve,值是 undefined——调用方拿到的本就是协调器返回的那个 promise。协议违规不会炸,只会安静地把整批调用变成 no-taskFn。

注 · 这套协议的契约分在两处读:类型面是 src/coordination/types.tsCoordinationContextCoordinator,运行期的强制(proceed/reuse 互斥、孤儿守卫、协调器同步抛错怎么收)是 src/coordination/coordinator-adapter.ts 里的 dispatchToCoordinator。想写自定义协调器,这两个文件就是全部输入。ctx.emit 看着像能发任意事件,但它的参数类型 CoordinationEventDraft 是从 coordination/events.ts 那个封闭 union 推导出来的,只认 reuse / queued / restart / orphaned-intent 四种 —— 第五种 type 过不了类型检查,自定义诊断得走 getEventBus(runtime) 自己发。

协调解决的是「同一件事被要求做了多次」。下一页处理另一半:一件事做了一次却没做成,重试与超时