系统设计 / 任务运行时 · 从一次 run() 到可观测的调度层 / 依赖与编排:从 after 到 pipeline 待审核 3 / 10
任务编排

依赖与编排:从 after 到 pipeline

单个任务的状态机讲完之后,剩下的问题是多个任务之间的先后。运行时给了两条路:一条是声明式的,提交任务时用 after 说明它等谁,图由运行时替调用方维护;另一条是命令式的,用 flow 把顺序写成一段程序,步骤之间的衔接由调用方控制。两条路解决的是同一件事,但可观测性与失败语义不同。本页先看依赖图,再看 flow 与它的三种并行步骤,最后落到一个具体的编排取舍:同一批数据过同一串阶段,逐阶段等齐与每个数据独立流过差多少。

1 · after 依赖图

RunOptions.after 收一个或一组前置条件,类型 TaskPrerequisite 的定义只有一行:任何带 promise 字段的对象。Task 句柄天然满足这个形状,一组 Task 组成的数组也满足。

after 的任务提交后不进队列,先停在 parked:deps。它不占并发槽,也不计入队列长度,全部依赖 resolve 之后才转 queued。任一依赖 reject 时,任务以 dep-failure 为原因取消,startedAt 始终为空。

依赖分 tracked 与 untracked 两类,判据是它是否为同一运行时拥有的 Task。tracked 依赖在 inspect(runtime).deps 上留下边,形状是「依赖方 id → 被依赖方 id 列表」;裸 { promise } 与外来运行时的 Task 属于 untracked,运行时只把它当作外部 future 等待,不记图边,并对它启用 dep-stuck 熔断。菱形依赖 A → B、A → C、(B, C) → D 的四条边全部是 tracked,可从 deps 表里逐条读出(lab/facts.test.ts §1)。

同一张菱形图上,三种故障位置给出三种结果(core/scenarios.test.ts §1):

  • 无故障:四个任务全部 success,D 最后完成。
  • A 失败:B、C、D 全部以 dep-failure 取消,三者的 taskFn 体一次都没执行。
  • 只有 C 失败:B 照常 success,D 因为两条入边里的一条断了而被取消。

取消沿依赖边传播,方向与任务与状态机里的父子级联无关。D 与 B、C 之间没有父子关系,它被取消只因为 after 数组里有一项 reject 了。

图 1-1 · 菱形依赖图在真实运行时上跑出的一次结果,结点色即当前状态,时间轴上未开跑的任务只留一段浅色等待条。可切换故障位置,观察取消沿依赖边波及到哪几个结点。

前置条件既然只要求一个 promise 字段,{ promise: task.started } 就同样合法,而它的语义与 after: task 完全不同:task.started 在依赖任务进入 running 时就 resolve,后继任务因此与依赖任务并发跑,只是起点被推后到依赖任务真正开跑的那一刻。实测一个 60ms 的主任务配一个这样挂上去的后继,事件序是「主任务开始 → 后继开始 → 主任务结束」(lab/facts.test.ts §1)。加载指示器只在主请求超过某个时长时才出现,用的就是这个偏移。

注 · 依赖图只管入队时机,不管数据传递。当一个任务的 taskFn 体开始执行时,它的全部 after 依赖保证已是 success 状态,直接读 dep.data 即可,无需再 await 一次。

2 · flow 与并行步骤

flow(runtime, ...steps) 把若干步骤按提交顺序串起来,每一步落定后下一步才被创建。flow 本身不是任务,不占队列与并发里说的那种并发槽,被展开出来的每个 step 才是任务。返回的 FlowHandle 有四个成员:promise 带末步的值,drained 在全部启动过的 taskFn 落定后 resolve 且永不 reject,correlationId 就是事件流里的 flowIdcancel() 取消这条 flow 派出的一切。

顺序之外的三种步骤都是 fan-out,即一步之内并发跑多个 taskFn。parallel(taskFns, options) 等全部 taskFn 落定,两个开关决定它的行为:

  • onError:默认 'stop',首个失败立即 abort 其余 taskFn 并让步骤 reject;'continue' 则等所有 taskFn 各自落定。
  • ordered:控制 onResult 的投递次序。默认按完成先后触发,置 true 后由一个重排缓冲区扣住提前到达的结果,直到前面的下标都投递完毕。

提交序 slow / fast / mid 的三个 taskFn,完成序必然是 fast → mid → slow。默认投递给出下标序列 1, 2, 0,ordered: true 给出 0, 1, 2,两次拿到的结果值本身相同(lab/facts.test.ts §2)。代价是快结果要等慢的前驱赶上来才放出去,只有下游逻辑依赖位置时才值得付。

onError 的两档在失败时分得很开。让 fast 抛错、其余两个各睡更久:'stop' 下三次回调一次成功也没有,下标 0 与 2 拿到的是 AbortError,步骤 promise 随之 reject;'continue' 下三次回调是「成功、失败、成功」,而步骤的 promise 仍然 resolve,失败只从 onResult 出来(lab/facts.test.ts §2)。这一条与仓库里 example/cases/ordered-parallel-stream.ts 的注释相左:那条注释说 parallel'continue' 下也会随任一 taskFn 失败而 reject,实测不成立。执行器只在 'stop' 分支上置停止标志并记下错误,'continue' 下池子抽干即 resolve。

图 2-1 · 一个 parallel 步骤里三个耗时不同的 taskFn,下方按 onResult 实际触发的先后排列。可切换 ordered 与 onError,观察回调次序以及失败时其余 taskFn 的结局。

另外两种步骤各有专攻。firstSuccess(taskFns) 以首个成功的值胜出并取消其余,全部失败时抛 AggregateError;它的惰性形态必须显式给 poolSize,胜者之后的元素根本不会被生产出来,适合镜像列表这类「有一个能用就够」的场合。flowDynamic(runtime, body) 把步骤写进一个 async body,await ctx.step(taskFn) 逐步派发,于是分支与循环都能写。同一个 body 里不允许并发调用 ctx.step(),要在 body 内部 fan-out 得写成 ctx.step(parallel([...]))

警示 · parallel 的数组形态与惰性形态在并发上限这一点上并不对称。数组重载的选项类型里根本没有 poolSize 字段,池子大小取运行时的 scheduler.concurrency;要按 poolSize 限流只能交给它一个生成器。六个 taskFn 提交到并发 6 的运行时,数组形态峰值并发是 6;同一批 taskFn 包成生成器再给 poolSize: 2,峰值是 2(lab/facts.test.ts §2)。

3 · pipeline 与 barrier

一批数据过同一串阶段,是编排里最常见的形状:六张图各要 fetch、parse、upload 三段。写法有两种,差别在于阶段之间要不要等齐。

逐阶段等齐的写法是 flow(rt, parallel(全部 fetch), parallel(全部 parse), parallel(全部 upload))。每个 parallel 步骤要等它的全部 taskFn 落定才交出控制权,阶段之间自然形成一道 barrier:任何一张图的 parse 都得等最慢的那张图 fetch 完。

pipeline(items, stages, options) 取消了这道屏障。它给每个 item 建一条独立的链,链上每一段的输入是上一段的输出,stage 签名是 (prev, item, index) => TaskFn,首段的 prev 就是 item 自身。poolSize 限的是同时在飞的 item 数,不是同时在飞的阶段数,所以一条链上的三段总是首尾相接。

同一批 6 个 item、pool 2、阶段耗时 60 / 20 / 40ms 时,两种写法的首个结果时刻可以直接算出来(core/scenarios.test.ts §3):pipeline 是 120ms,即一条链三段之和;barrier 是 280ms,即 3 轮 fetch 加 3 轮 parse 再加一次 upload;barrier 的总墙钟是 360ms。实跑同一组场景,pipeline 的首个结果确实早于 barrier,两者最终都把 6 个 item 送进了 upload。这个差距随 item 数与阶段耗时的不均衡一起放大,而总墙钟未必改善——barrier 浪费的是交付时机,不一定是吞吐。

图 3-1 · 六个 item 走 fetch → parse → upload 三段,一行一个 item,虚线标出首个结果时刻。可调 poolSize 后分别运行两种模式,观察阶段条是连成阶梯还是分层成三块。

pipeline 的阶段是内联在链的 ctx 上执行的,并不派生子任务。源码注释给出了理由:子任务会在父任务等待期间额外占住一个并发槽,poolSize 一旦逼近运行时并发就会死锁。代价是 devtools 里看不到「一个 stage 一个任务」,链改为在每段之前发一条 ctx.progress 事件补偿,stageLabels 就写进那条事件的 message。同样因为阶段是内联的,pause(runtime) 的粒度落在链的边界:在飞的链会把剩余阶段一口气跑完,被挡住的只有还没开始的 item。

警示 · pipeline 的默认失败策略与同目录的 parallel 不一样,签名上看不出来,得读 PipelineOptions 的默认值。默认是 onError: 'drop':抛错的 item 跳过自己剩下的阶段,同批其他 item 照常走完,失败以 { ok: false } 的形式从 onResult 出来,步骤的 promise 不 reject。四个 item、第 2 个在首段抛错时,onResult 依次给出成功、失败、成功、成功,另外三个 item 全部走到了末段(lab/facts.test.ts §3)。要让首个失败掐掉整批,得显式写 onError: 'stop'

依赖图与 flow 管的都是「谁在谁之后」,回答不了「同一件事被连着调用了五次该执行几次」。下一页是时间维度上的收敛:时间窗口 shaper