如何用拓扑排序实现带依赖的嵌套Job任务并行执行?
基于JGraphT实现带并行执行的Job依赖调度
问题分析
TopologicalOrderIterator给出的是线性拓扑顺序,但你的场景需要同一依赖层级的Job并行执行:JOB1和JOB3无前置依赖需同时启动;JOB2依赖JOB1,需等JOB1完成后启动;JOB_BOX需等所有子Job完成后标记完成。线性拓扑排序无法直接满足并行需求,需要基于入度的分层调度方案。
可行方案:基于入度的分层并行调度
你提出的入度判断方案完全可行,本质是Kahn拓扑排序算法的扩展。核心逻辑是:
- 跟踪每个Job的入度(未完成的前置依赖数量)
- 每次执行所有入度为0的Job(这些Job无未完成依赖,可并行)
- 单个Job完成后,更新其下游Job的入度,当下游Job入度变为0时,加入下一批执行队列
- 重复上述步骤直到所有Job执行完毕
具体实现代码
针对你的场景,用JGraphT实现该逻辑:
import org.jgrapht.graph.DirectedAcyclicGraph; import org.jgrapht.graph.DefaultEdge; import java.util.*; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; public class JobScheduler { public static void main(String[] args) throws InterruptedException { DirectedAcyclicGraph<String, DefaultEdge> jobGraph = new DirectedAcyclicGraph<>(DefaultEdge.class); jobGraph.addVertex("JOB_BOX"); jobGraph.addVertex("JOB1"); jobGraph.addVertex("JOB2"); jobGraph.addVertex("JOB3"); // 边规则:前置Job -> 依赖它的Job jobGraph.addEdge("JOB1", "JOB2"); jobGraph.addEdge("JOB1", "JOB_BOX"); jobGraph.addEdge("JOB2", "JOB_BOX"); jobGraph.addEdge("JOB3", "JOB_BOX"); // 初始化每个节点的入度 Map<String, Integer> inDegreeMap = new HashMap<>(); for (String vertex : jobGraph.vertexSet()) { inDegreeMap.put(vertex, jobGraph.inDegreeOf(vertex)); } // 线程池处理并行Job ExecutorService executor = Executors.newFixedThreadPool(3); // 待执行队列:入度为0的节点 Queue<String> queue = new LinkedList<>(); inDegreeMap.entrySet().stream() .filter(entry -> entry.getValue() == 0) .forEach(entry -> queue.add(entry.getKey())); while (!queue.isEmpty()) { // 取出当前所有可并行执行的Job List<String> currentBatch = new ArrayList<>(queue); queue.clear(); // 提交当前批次Job到线程池 for (String job : currentBatch) { executor.submit(() -> { System.out.println("开始执行: " + job); // 模拟Job执行耗时 try { TimeUnit.SECONDS.sleep(1); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("完成执行: " + job); // 同步更新下游Job的入度 synchronized (inDegreeMap) { for (DefaultEdge edge : jobGraph.outgoingEdgesOf(job)) { String nextJob = jobGraph.getEdgeTarget(edge); int newDegree = inDegreeMap.get(nextJob) - 1; inDegreeMap.put(nextJob, newDegree); // 入度为0则加入下一批队列 if (newDegree == 0) { queue.add(nextJob); } } } }); } // 等待当前批次所有Job完成,再处理下一批 executor.shutdown(); executor.awaitTermination(5, TimeUnit.SECONDS); executor = Executors.newFixedThreadPool(3); } executor.shutdown(); System.out.println("所有Job执行完成,JOB_BOX标记为已完成"); } }
代码说明
- 入度跟踪:用
inDegreeMap记录每个Job未完成的前置依赖数,初始值为节点的入度 - 批次并行:每次取出所有入度为0的Job组成批次,提交到线程池并行执行
- 依赖更新:Job完成后同步更新下游Job的入度,入度归0的Job进入下一批执行队列
- 批次等待:每批执行完成后等待所有线程结束,确保依赖顺序不被破坏
嵌套Job Box扩展
若需支持嵌套Job Box,只需将嵌套Box当作普通Job节点处理:
- 嵌套Box的入度依赖于其内部所有子Job的完成
- 内部子Job的依赖关系用相同图结构表示,递归应用上述调度逻辑即可
参考思路(原Stack Overflow内容翻译)
原参考内容的核心逻辑:
- 拓扑排序的核心是处理依赖,而并行执行需要识别无依赖节点组
- Kahn算法(基于入度的拓扑排序)天然适配分层并行场景,可按层级处理节点
- 实现时需维护节点入度,节点完成后更新下游入度,入度为0时加入执行队列
内容的提问来源于stack exchange,提问作者vkp
相关产品推荐
相关产品推荐

