Skip to content

依赖任务怎样安全地跑起来:拓扑排序只是第一步 ​

报名流程里的「名额预留」第一次调用已经扣了名额,只是响应超时;执行器重试一次,同一个报名占掉了 2 个名额。算出执行顺序之后,还有配置校验、并发、失败传播和重试这几件事要处理。

一个活动报名要经过几步:先做资格校验,再预留名额,然后费用计算、通知准备、风控复核三件事可以同时做,最后提交。步骤之间的先后关系写在配置里,由一个执行器按依赖关系调度。

把这张依赖图排成一个顺序,用 Kahn 算法几十行就够了,工程里用得上的十个算法套路 里有完整代码。执行器要做的事比这多:坏配置要在发布前拒绝;两次执行要给出相同的顺序;能并行的步骤要并行,但同时运行的数量要有上限;一个步骤失败后,要决定哪些步骤还能继续;重试不能让副作用重复生效。下面每一节对应其中一件事,结果都来自同一个实验程序(JDK 21.0.5),见文末配套实验。

一、发布前就拒绝坏配置 ​

配置出错有四种情况:两个步骤互相依赖(环)、依赖了不存在的步骤(缺失)、同一个步骤定义了两次(重复),以及依赖了环上步骤的下游。前三种是错误本身,第四种只是被错误卡住。

资格校验名额预留真实的环费用计算受环阻塞提交受环阻塞发票配置里没有这个节点缺失节点通知准备第 1 次定义通知准备第 2 次定义重复节点四类错误在发布前一次报全,运行时不会出现「一直等不到依赖」的节点
图 1 · 一份有问题的报名流程配置(箭头 a → b 表示 a 完成后 b 才能运行):Kahn 算法结束后剩下 4 个节点,只有资格校验和名额预留在环上;费用计算和提交只是被环卡住,缺失的发票和重复的通知准备要单独报告

Kahn 算法结束后,剩下的节点要么在环上,要么依赖了环上的节点。把剩下的节点一股脑报成「循环依赖」,配置人员会去改一条本身没有问题的规则。报告真实的环有一个简单的做法:剩下的每个节点至少有一个依赖也在剩下的集合里,从任意一个出发沿依赖往回走,第一次走到重复节点时,截出来的那一段就是环。

java
// blocked:Kahn 算法结束后入度仍大于 0 的节点
Set<String> blocked = ...;
List<String> path = new ArrayList<>();
Map<String, Integer> seenAt = new HashMap<>();
String cur = blocked.iterator().next();
while (!seenAt.containsKey(cur)) {
    seenAt.put(cur, path.size());
    path.add(cur);
    cur = byId.get(cur).deps().stream()
            .filter(blocked::contains).sorted().findFirst().orElseThrow();
}
List<String> cycle = new ArrayList<>(path.subList(seenAt.get(cur), path.size()));
cycle.add(cur);   // 其余的 blocked 节点单独报「受环阻塞」

对图 1 的配置,校验一次报出四条(这里的箭头表示「依赖」):

text
重复节点: 通知准备;缺失节点: 提交 依赖的 发票 不存在;
循环依赖: 名额预留 -> 资格校验 -> 名额预留;受环阻塞: [提交, 费用计算]

这四类检查都放在发布配置的时候做。运行时才发现环,执行器会停在一个永远等不到依赖的步骤上,没有报错,也没有超时。

二、同一张图,两次排出的顺序不一样 ​

拓扑顺序通常不唯一。费用计算、通知准备、风控复核之间没有依赖,谁先谁后都合法。问题在于,普通队列的出队顺序取决于节点是按什么顺序加进来的,而这又取决于配置的声明顺序:

ready 集合按原顺序声明倒序声明
ArrayDeque(先进先出)费用计算、通知准备、风控复核风控复核、通知准备、费用计算
按名称排序(PriorityQueue、TreeSet)费用计算、通知准备、风控复核费用计算、通知准备、风控复核

如果执行日志要逐次对比、测试要断言顺序,或者某个步骤的输出会影响下一次的缓存键,就用一个带固定排序键的集合保存 ready 节点。排序键可以是名称,也可以是配置里显式写的优先级。顺序是否稳定,要在设计执行器时就定下来,不要依赖 HashMap 碰巧的遍历顺序。

三、并发执行:关键路径是下限,并发上限决定能不能接近它 ​

每个步骤完成后,把它的后继入度减一,减到零的放进 ready 集合;只要正在运行的数量没到上限,就从 ready 集合里取下一个提交给线程池。核心循环如下(done 是 ExecutorCompletionService,limit 是并发上限):

java
while (!ready.isEmpty() || !inFlight.isEmpty()) {
    while (!ready.isEmpty() && inFlight.size() < limit) {
        String id = ready.pollFirst();
        inFlight.put(id, done.submit(() -> { run(id); return id; }));
    }
    Future<String> f = done.take();                 // 等任意一个步骤完成
    String id = idOf(f);
    inFlight.remove(id);
    try {
        f.get();
        for (String n : dependents.getOrDefault(id, List.of())) {
            if (pending.merge(n, -1, Integer::sum) == 0) ready.add(n);
        }
    } catch (ExecutionException e) {
        onFailure(id, e.getCause());                                    // 见第四节
    }
}

实验里资格校验和名额预留各 100ms,三个并行分支各 200ms,提交 50ms。所有步骤耗时加起来是 850ms,关键路径(资格校验 → 名额预留 → 任一分支 → 提交)是 450ms:

并发上限总耗时同时运行的步骤数(最大)
1881ms1
2678ms2
3471ms3

上限为 1 时退化成顺序执行。上限为 2 时,三个分支要分两批跑,比关键路径多出一个分支的时间。上限为 3 时接近关键路径,再往上加也不会更快。所以上限不是越大越好:它要和下游能承受的并发一起定,道理和 线程池参数怎么定 相同。「同时运行的步骤数」要作为指标记录下来,它是判断上限是否生效的直接证据。

四、一个步骤失败之后 ​

费用计算失败时,提交显然不能再执行,因为它依赖费用计算。有争议的是另外两个分支:通知准备和风控复核和费用计算无关,还在运行。

策略一:跳过失败节点的后继资格校验、名额预留成功费用计算失败通知准备成功风控复核成功提交跳过策略二:立即取消整张图资格校验、名额预留成功费用计算失败通知准备取消风控复核取消提交跳过350ms 费用计算失败0ms100200300400450
图 2 · 费用计算在 350ms 失败(时长为实验中的设定值):「跳过后继」让无关的通知准备和风控复核跑完,只跳过提交;「立即取消」中断仍在运行的分支。两种策略下提交都不会执行

两种策略的实测结果:

策略费用计算通知准备、风控复核提交
跳过失败步骤的后继失败成功跳过(依赖费用计算失败)
立即取消整张图失败取消跳过

选哪一种由业务决定,不由执行器决定。通知准备只是生成文案,跑完了也没有副作用,「跳过后继」可以少做一次重跑。如果某个分支会对外产生副作用(发短信、调用第三方),而整个报名已经注定失败,就应该「立即取消」,并且在取消时考虑已经完成的步骤要不要补偿。

取消依赖任务响应中断。实验里的步骤在 Thread.sleep 中等待,Future.cancel(true) 能让它立即停下;如果步骤阻塞在不响应中断的调用上,「立即取消」只能做到不再启动新的步骤。

每个步骤最终都要落到一个明确的状态上:成功、失败、跳过或取消,并记录原因(「依赖 费用计算 失败」)。排查时,跳过和失败的区别很重要:跳过的步骤本身没有问题,不需要去看它的日志。

五、重试之前先有幂等键 ​

名额预留调用库存服务,库存服务已经扣减了名额,但响应在网络上超时了。执行器看到的是一次失败,于是重试。

调用方式调用次数实际占用名额
不带幂等键22
带幂等键(执行编号 + 步骤名)21

执行器能做的是为每次执行生成一个唯一的执行编号(run_id),把「执行编号 + 步骤名」作为幂等键传给下游,同一个步骤的每次重试都带同一个键。去重发生在被调用的服务里:它要在同一个事务里记录这个键并执行扣减,重复的键直接返回上一次的结果。具体的做法和并发边界见 可复用的服务端组件 中的幂等键实验,订单场景见 订单、库存与数据一致性。

幂等键不能用重试次数、时间戳或者执行器所在的机器名来拼,这些在重试时都会变。为什么超时之后不能直接重试、哪些失败可以放心重试,见 一次 HTTP 调用超时了,对方到底执行了没有。

六、什么时候不必自己写 ​

步骤少而且固定时,直接在代码里调用更清楚:几个 CompletableFuture 加一次 allOf 就能表达「三个分支并行、全部完成后提交」,失败处理也写在调用处,读代码的人不需要先理解一个执行器。

需要自己写执行器,通常是因为步骤由配置决定、会随业务变化,而且需要统一的校验、状态记录和失败策略。如果还需要在进程崩溃后从中断的步骤继续、需要人工审批这类跨天的等待,就已经是工作流编排的范围了,应该先评估现成的工作流引擎,而不是在这个执行器上继续加功能。单个业务对象的状态迁移(报名从待确认到已确认)是另一个问题,见 状态机与工作流。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权