You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

带依赖的任务调度与Worker分配实现技术问询

实现满足约束的任务调度Worker分配方案

需求与约束

我要实现一个简单的任务调度系统,给定一组任务(每个任务已知耗时effort,部分有前置依赖)和指定数量的独立Worker,需要计算任务分配方案,满足:

  1. 任务必须在所有依赖任务完成后才能启动
  2. 任何时刻不能同时存在空闲Worker和可执行任务(无未完成依赖的待执行任务)

示例

给定任务集合:

  • Task A:effort=1,无依赖
  • Task B:effort=2,无依赖
  • Task C:effort=1,依赖Task A、B

1个Worker时的执行方案

Worker IDTaskStart TimeEnd Time
1Task A01
1Task B13
1Task C34

2个Worker时的执行方案

Worker IDTaskStart TimeEnd Time
1Task A01
2Task B02
1Task C23

*注: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; }
    }
}

关键逻辑说明

  1. 拓扑排序的作用:确保任务处理顺序符合依赖关系,不会出现先处理依赖任务的情况
  2. Worker堆的作用:每次选择最早空闲的Worker,保证资源利用率最大化
  3. 就绪任务池的维护:任务完成后立即检查下游任务的依赖状态,满足条件的任务及时加入就绪池,避免空闲Worker等待
  4. 任务启动时间计算:取Worker空闲时间和任务依赖最晚完成时间的最大值,严格满足依赖约束

约束验证

  • 约束1:任务启动时间取依赖最晚完成时间,确保所有依赖完成后才启动
  • 约束2:只要有就绪任务,最早空闲的Worker会立即分配任务,不会出现空闲Worker与就绪任务并存的情况

内容的提问来源于stack exchange,提问作者Barcelona

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 22:50:40