带依赖的任务调度与Worker分配实现技术问询
实现满足约束的任务调度Worker分配方案
需求与约束
我要实现一个简单的任务调度系统,给定一组任务(每个任务已知耗时effort,部分有前置依赖)和指定数量的独立Worker,需要计算任务分配方案,满足:
- 任务必须在所有依赖任务完成后才能启动
- 任何时刻不能同时存在空闲Worker和可执行任务(无未完成依赖的待执行任务)
示例
给定任务集合:
- Task A:effort=1,无依赖
- Task B:effort=2,无依赖
- Task C:effort=1,依赖Task A、B
1个Worker时的执行方案
| Worker ID | Task | Start Time | End Time |
|---|---|---|---|
| 1 | Task A | 0 | 1 |
| 1 | Task B | 1 | 3 |
| 1 | Task C | 3 | 4 |
2个Worker时的执行方案
| Worker ID | Task | Start Time | End Time |
|---|---|---|---|
| 1 | Task A | 0 | 1 |
| 2 | Task B | 0 | 2 |
| 1 | Task C | 2 | 3 |
*注:Worker完成任务后可在同一时间步启动新任务;Worker与任务的匹配顺序不做要求,只要满足约束即可。
已定义的API契约
public interface Task extends Iterable<Task> { String getName(); int getEffort(); List<Task> getDependencies(); } public interface TaskExecution { int getEnd(); Task getTask(); int getWorkerId(); int getStart(); } public interface Scheduler { List<TaskExecution> estimate(Iterable<Task> tasks, int numWorkers); }
当前进展与疑问
我已经通过图结构和拓扑排序完成了任务的依赖排序,但不知道如何实现Worker分配以最大化并行度。现有Graph类代码如下:
public class Graph<T> { static class Node<E> { E source; E destination; Node(E source, E destination) { this.source = source; this.destination = destination; } } private int vertices; private LinkedList<Node<T>>[] adjList; public void addEgde(T source, T destination, List<T> list) { Node node = new Node(source, destination); adjList[list.indexOf(source)].addFirst(node); } public void topologicalSorting(List<T> list) { // 已实现拓扑排序逻辑 } }
请问如何实现满足要求的Worker分配,以实现最大程度的并行执行?
解决方案:基于事件驱动的Worker调度算法
要实现最大化并行度的Worker分配,核心是跟踪每个Worker的空闲时间,同时动态维护可执行任务池,确保一旦有Worker空闲或有任务满足依赖条件,就立即分配任务。以下是具体实现步骤:
1. 预处理任务依赖与状态
首先,基于拓扑排序后的任务列表,我们需要维护以下数据结构:
Map<Task, List<Task>> taskToDownstream:记录每个任务的下游任务(依赖当前任务的任务),用于任务完成后更新下游状态Map<Task, Integer> dependencyCount:记录每个任务剩余未完成的依赖数量Map<Task, Integer> taskLatestDependencyEndTime:记录每个任务所有依赖的最晚完成时间(决定任务最早启动时间)Map<Task, Integer> taskCompletedTime:记录已完成任务的结束时间PriorityQueue<Task> readyTasks:存储当前就绪(所有依赖已完成)的待执行任务,可按任务耗时排序(短任务优先,提升资源利用率)
2. 跟踪Worker的空闲状态
使用**最小堆(优先队列)**管理Worker的空闲时间,堆中元素包含Worker ID和其空闲时间,这样能快速获取最早空闲的Worker,保证资源不被浪费。
3. 事件驱动的调度循环
调度流程的核心是循环处理两种场景:
- Worker空闲时,立即从就绪任务池中分配任务
- 任务完成后,更新下游任务的依赖状态,将满足条件的下游任务加入就绪池
完整实现代码
public class SchedulerImpl implements Scheduler { @Override public List<TaskExecution> estimate(Iterable<Task> tasks, int numWorkers) { List<TaskExecution> executions = new ArrayList<>(); // 1. 获取拓扑排序后的任务列表 List<Task> topoSortedTasks = getTopologicallySortedTasks(tasks); // 2. 初始化依赖与状态映射 Map<Task, List<Task>> taskToDownstream = new HashMap<>(); Map<Task, Integer> dependencyCount = new HashMap<>(); Map<Task, Integer> taskLatestDependencyEndTime = new HashMap<>(); Map<Task, Integer> taskCompletedTime = new HashMap<>(); for (Task task : topoSortedTasks) { dependencyCount.put(task, task.getDependencies().size()); taskLatestDependencyEndTime.put(task, 0); // 构建下游任务映射 for (Task dep : task.getDependencies()) { taskToDownstream.computeIfAbsent(dep, k -> new ArrayList<>()).add(task); } } // 3. 初始化就绪任务池(短任务优先) PriorityQueue<Task> readyTasks = new PriorityQueue<>(Comparator.comparingInt(Task::getEffort)); for (Task task : topoSortedTasks) { if (task.getDependencies().isEmpty()) { readyTasks.add(task); } } // 4. 初始化Worker空闲堆:按空闲时间升序排列 PriorityQueue<WorkerIdle> workerHeap = new PriorityQueue<>(Comparator.comparingInt(WorkerIdle::getIdleTime)); for (int i = 1; i <= numWorkers; i++) { workerHeap.add(new WorkerIdle(0, i)); } // 5. 循环调度直到所有任务完成 while (taskCompletedTime.size() < topoSortedTasks.size()) { WorkerIdle earliestIdleWorker = workerHeap.poll(); int workerId = earliestIdleWorker.getWorkerId(); int currentWorkerIdleTime = earliestIdleWorker.getIdleTime(); if (!readyTasks.isEmpty()) { // 分配就绪任务 Task task = readyTasks.poll(); // 任务最早启动时间:取Worker空闲时间和依赖最晚完成时间的最大值 int startTime = Math.max(currentWorkerIdleTime, taskLatestDependencyEndTime.get(task)); int endTime = startTime + task.getEffort(); // 记录执行信息 executions.add(new TaskExecutionImpl(task, workerId, startTime, endTime)); taskCompletedTime.put(task, endTime); // 更新下游任务状态 List<Task> downstreamTasks = taskToDownstream.getOrDefault(task, Collections.emptyList()); for (Task downstream : downstreamTasks) { dependencyCount.put(downstream, dependencyCount.get(downstream) - 1); // 更新下游任务的最晚依赖完成时间 int latestEnd = Math.max(taskLatestDependencyEndTime.get(downstream), endTime); taskLatestDependencyEndTime.put(downstream, latestEnd); // 若所有依赖完成,加入就绪池 if (dependencyCount.get(downstream) == 0) { readyTasks.add(downstream); } } // 将Worker重新加入堆,更新空闲时间 workerHeap.add(new WorkerIdle(endTime, workerId)); } else { // 无就绪任务,等待到下一个任务完成时间 int nextCompletionTime = Collections.min(taskCompletedTime.values()); workerHeap.add(new WorkerIdle(nextCompletionTime, workerId)); } } return executions; } // 调用已实现的拓扑排序方法,返回排序后的任务列表 private List<Task> getTopologicallySortedTasks(Iterable<Task> tasks) { List<Task> taskList = new ArrayList<>(); for (Task task : tasks) { taskList.add(task); } Graph<Task> graph = new Graph<>(); // 构建依赖边:依赖任务 -> 当前任务 for (Task task : taskList) { for (Task dep : task.getDependencies()) { graph.addEgde(dep, task, taskList); } } graph.topologicalSorting(taskList); // 假设topologicalSorting方法会修改taskList为排序后的结果,可根据实际实现调整 return taskList; } // 内部类:存储Worker空闲时间与ID private static class WorkerIdle { private int idleTime; private int workerId; public WorkerIdle(int idleTime, int workerId) { this.idleTime = idleTime; this.workerId = workerId; } public int getIdleTime() { return idleTime; } public int getWorkerId() { return workerId; } } // TaskExecution接口实现类 private static class TaskExecutionImpl implements TaskExecution { private Task task; private int workerId; private int start; private int end; public TaskExecutionImpl(Task task, int workerId, int start, int end) { this.task = task; this.workerId = workerId; this.start = start; this.end = end; } @Override public int getEnd() { return end; } @Override public Task getTask() { return task; } @Override public int getWorkerId() { return workerId; } @Override public int getStart() { return start; } } }
关键逻辑说明
- 拓扑排序的作用:确保任务处理顺序符合依赖关系,不会出现先处理依赖任务的情况
- Worker堆的作用:每次选择最早空闲的Worker,保证资源利用率最大化
- 就绪任务池的维护:任务完成后立即检查下游任务的依赖状态,满足条件的任务及时加入就绪池,避免空闲Worker等待
- 任务启动时间计算:取Worker空闲时间和任务依赖最晚完成时间的最大值,严格满足依赖约束
约束验证
- 约束1:任务启动时间取依赖最晚完成时间,确保所有依赖完成后才启动
- 约束2:只要有就绪任务,最早空闲的Worker会立即分配任务,不会出现空闲Worker与就绪任务并存的情况
内容的提问来源于stack exchange,提问作者Barcelona
相关产品推荐
相关产品推荐

