依赖任务怎样安全地跑起来:拓扑排序只是第一步
报名流程里的「名额预留」第一次调用已经扣了名额,只是响应超时;执行器重试一次,同一个报名占掉了 2 个名额。算出执行顺序之后,还有配置校验、并发、失败传播和重试这几件事要处理。
一个活动报名要经过几步:先做资格校验,再预留名额,然后费用计算、通知准备、风控复核三件事可以同时做,最后提交。步骤之间的先后关系写在配置里,由一个执行器按依赖关系调度。
把这张依赖图排成一个顺序,用 Kahn 算法几十行就够了,工程里用得上的十个算法套路 里有完整代码。执行器要做的事比这多:坏配置要在发布前拒绝;两次执行要给出相同的顺序;能并行的步骤要并行,但同时运行的数量要有上限;一个步骤失败后,要决定哪些步骤还能继续;重试不能让副作用重复生效。下面每一节对应其中一件事,结果都来自同一个实验程序(JDK 21.0.5),见文末配套实验。
一、发布前就拒绝坏配置
配置出错有四种情况:两个步骤互相依赖(环)、依赖了不存在的步骤(缺失)、同一个步骤定义了两次(重复),以及依赖了环上步骤的下游。前三种是错误本身,第四种只是被错误卡住。
Kahn 算法结束后,剩下的节点要么在环上,要么依赖了环上的节点。把剩下的节点一股脑报成「循环依赖」,配置人员会去改一条本身没有问题的规则。报告真实的环有一个简单的做法:剩下的每个节点至少有一个依赖也在剩下的集合里,从任意一个出发沿依赖往回走,第一次走到重复节点时,截出来的那一段就是环。
// 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 的配置,校验一次报出四条(这里的箭头表示「依赖」):
重复节点: 通知准备;缺失节点: 提交 依赖的 发票 不存在;
循环依赖: 名额预留 -> 资格校验 -> 名额预留;受环阻塞: [提交, 费用计算]这四类检查都放在发布配置的时候做。运行时才发现环,执行器会停在一个永远等不到依赖的步骤上,没有报错,也没有超时。
二、同一张图,两次排出的顺序不一样
拓扑顺序通常不唯一。费用计算、通知准备、风控复核之间没有依赖,谁先谁后都合法。问题在于,普通队列的出队顺序取决于节点是按什么顺序加进来的,而这又取决于配置的声明顺序:
| ready 集合 | 按原顺序声明 | 倒序声明 |
|---|---|---|
ArrayDeque(先进先出) | 费用计算、通知准备、风控复核 | 风控复核、通知准备、费用计算 |
按名称排序(PriorityQueue、TreeSet) | 费用计算、通知准备、风控复核 | 费用计算、通知准备、风控复核 |
如果执行日志要逐次对比、测试要断言顺序,或者某个步骤的输出会影响下一次的缓存键,就用一个带固定排序键的集合保存 ready 节点。排序键可以是名称,也可以是配置里显式写的优先级。顺序是否稳定,要在设计执行器时就定下来,不要依赖 HashMap 碰巧的遍历顺序。
三、并发执行:关键路径是下限,并发上限决定能不能接近它
每个步骤完成后,把它的后继入度减一,减到零的放进 ready 集合;只要正在运行的数量没到上限,就从 ready 集合里取下一个提交给线程池。核心循环如下(done 是 ExecutorCompletionService,limit 是并发上限):
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:
| 并发上限 | 总耗时 | 同时运行的步骤数(最大) |
|---|---|---|
| 1 | 881ms | 1 |
| 2 | 678ms | 2 |
| 3 | 471ms | 3 |
上限为 1 时退化成顺序执行。上限为 2 时,三个分支要分两批跑,比关键路径多出一个分支的时间。上限为 3 时接近关键路径,再往上加也不会更快。所以上限不是越大越好:它要和下游能承受的并发一起定,道理和 线程池参数怎么定 相同。「同时运行的步骤数」要作为指标记录下来,它是判断上限是否生效的直接证据。
四、一个步骤失败之后
费用计算失败时,提交显然不能再执行,因为它依赖费用计算。有争议的是另外两个分支:通知准备和风控复核和费用计算无关,还在运行。
两种策略的实测结果:
| 策略 | 费用计算 | 通知准备、风控复核 | 提交 |
|---|---|---|---|
| 跳过失败步骤的后继 | 失败 | 成功 | 跳过(依赖费用计算失败) |
| 立即取消整张图 | 失败 | 取消 | 跳过 |
选哪一种由业务决定,不由执行器决定。通知准备只是生成文案,跑完了也没有副作用,「跳过后继」可以少做一次重跑。如果某个分支会对外产生副作用(发短信、调用第三方),而整个报名已经注定失败,就应该「立即取消」,并且在取消时考虑已经完成的步骤要不要补偿。
取消依赖任务响应中断。实验里的步骤在 Thread.sleep 中等待,Future.cancel(true) 能让它立即停下;如果步骤阻塞在不响应中断的调用上,「立即取消」只能做到不再启动新的步骤。
每个步骤最终都要落到一个明确的状态上:成功、失败、跳过或取消,并记录原因(「依赖 费用计算 失败」)。排查时,跳过和失败的区别很重要:跳过的步骤本身没有问题,不需要去看它的日志。
五、重试之前先有幂等键
名额预留调用库存服务,库存服务已经扣减了名额,但响应在网络上超时了。执行器看到的是一次失败,于是重试。
| 调用方式 | 调用次数 | 实际占用名额 |
|---|---|---|
| 不带幂等键 | 2 | 2 |
| 带幂等键(执行编号 + 步骤名) | 2 | 1 |
执行器能做的是为每次执行生成一个唯一的执行编号(run_id),把「执行编号 + 步骤名」作为幂等键传给下游,同一个步骤的每次重试都带同一个键。去重发生在被调用的服务里:它要在同一个事务里记录这个键并执行扣减,重复的键直接返回上一次的结果。具体的做法和并发边界见 可复用的服务端组件 中的幂等键实验,订单场景见 订单、库存与数据一致性。
幂等键不能用重试次数、时间戳或者执行器所在的机器名来拼,这些在重试时都会变。为什么超时之后不能直接重试、哪些失败可以放心重试,见 一次 HTTP 调用超时了,对方到底执行了没有。
六、什么时候不必自己写
步骤少而且固定时,直接在代码里调用更清楚:几个 CompletableFuture 加一次 allOf 就能表达「三个分支并行、全部完成后提交」,失败处理也写在调用处,读代码的人不需要先理解一个执行器。
需要自己写执行器,通常是因为步骤由配置决定、会随业务变化,而且需要统一的校验、状态记录和失败策略。如果还需要在进程崩溃后从中断的步骤继续、需要人工审批这类跨天的等待,就已经是工作流编排的范围了,应该先评估现成的工作流引擎,而不是在这个执行器上继续加功能。单个业务对象的状态迁移(报名从待确认到已确认)是另一个问题,见 状态机与工作流。
配套实验
- codesphere-labs/system-design/dependency-graph-execution:配置校验、稳定顺序、并发上限 1—3 的调度、两种失败策略、带与不带幂等键的重试(验证记录)
参考资料
- Arthur B. Kahn,Topological sorting of large networks,Communications of the ACM,1962
- JDK 21 API:ExecutorCompletionService
- JDK 21 API:Future.cancel