Task Scheduling Algorithms
This document is a work in progress and will be updated regularly.
Introduction
Scheduling a task graph on networked computers is a fundamental problem in distributed computing. Essentially, the goal is to assign computational tasks to different compute nodes in such a way that minimizes/maximizes some performance metric (e.g., total execution time, energy consumption, throughput, etc.).
We will focus on the task scheduling problem concerning heterogeneous task graphs and compute networks with the objective of minimizing makespan (total execution time) under the related machines model1.
It is common to model distributed applications as task graphs, where nodes represent computational tasks and directed edges represent precedence constraints and the flow of input/output data. As a result, task scheduling pops up all over the place - from machine learning and scientific workflows, to IoT/edge computing applications, to data processing pipelines used all over industry.
Figure 1 depicts a scientific workflow application used by Caltech astronomers to generate science-grade mosaics from astronomical imagery[montage].
Figure 1a: Montage astronomical image
Figure 1b: Montage scientific workflow structure
Problem Definition
Let us denote the task graph as , where is the set of tasks and contains the directed edges or dependencies between these tasks. An edge implies that the output from task is required input for task . Thus, task cannot start executing until it has received the output of task . This is often referred to as a precedence constraint.
For a given task , its compute cost is represented by and the size of the data exchanged between two dependent tasks, , is .
Let denote the compute node network, where is a complete undirected graph. is the set of nodes and is the set of edges. The compute speed of a node is and the communication strength between nodes is .
Under the related machines model[graham], the execution time of a task on a node is , and the data communication time between tasks from node to node (i.e., executes on and executes on ) is .
The goal is to schedule the tasks on different compute nodes in such a way that minimizes the makespan (total execution time) of the task graph.
Let denote a task scheduling algorithm. Given a problem instance which represents a network/task graph pair, let denote the schedule produced by for . A schedule is a mapping from each task to a triple where is the node on which the task is scheduled, is the start time, and is the end time.
A valid schedule must satisfy the following properties:
-
All tasks must be scheduled: for all , must exist such that and .
-
All tasks must have valid start and end times:
-
Only one task can be scheduled on a node at a time (i.e., their start/end times cannot overlap):
-
A task cannot start executing until all of its dependencies have finished executing and their outputs have been received at the node on which the task is scheduled:
Figure 2a: Example task graph
Figure 2b: Example compute network
Figure 2c: Example schedule (Gantt chart)
We define the makespan of the schedule as the time at which the last task finishes executing:
Example 1
Take a look at the task graph, network, and schedule in Figure 2. Let us start by verifying that this is a valid schedule for the problem instance (network/task graph pair).
First, task is scheduled to run on node . Clearly this is valid, since has no dependencies. When finishes running at time , which is valid since the cost of task is and the speed of node is ().
Then, immediately starts running at time on node . Again, this is clearly valid since there is no communication delay in sending the outputs from task to another node before running task .
Task , on the other hand, is scheduled to run on node . In this case, unit of output data from task must be sent to node as input data to task . The communication link between nodes and is , so this communication takes units of time. Thus, the start time of task is valid since it is exactly units of time after task terminates.
It's easy to verify that tasks and have valid runtimes according to their costs and the speeds of the nodes they're running on.
Finally, task is scheduled to run on node . Before it can start running, though, the units of output data from task must be sent from node to node over a communication link of strength . Thus, the start time of task is correct ( units of time after task 's finish time).
Thus, the schedule in Figure 2c is valid and has a makespan of .
The HEFT Scheduling Algorithm
This task scheduling problem has long been known to be NP-Hard and was recently shown to also be not polynomial-time approximable within a constant factor[inapproximable]. As a result, many heuristic algorithms that aren't guaranteed to produce an optimal schedule but that, in practice, have been shown to work reasonably well have been proposed over the past decades.
One of the most commonly used of these algorithms is HEFT (Heterogeneous Earliest Finish Time)[heft]. HEFT is a list-scheduling algorithm, which essentially means it first computes priorities for each of the tasks in the task graph and then schedules the tasks greedily in order of their priority on the "best" node (the one that minimizes the task's finish time, given previously scheduled tasks).
Here is a summary of the algorithm:
-
Calculate average compute times for each task:
-
Calculate average communication times for each dependency:
-
Calculate the upward rank of each task (recursively):
-
In descending order of task upward ranks, greedily schedule each task on the node that minimizes its earliest possible finish time given previously scheduled tasks.
HEFT Example Calculations
| Task | |
|---|---|
| 2/3 | |
| 2 | |
| 4/3 | |
| 2/3 |
Table 1: Average compute times for each task
| 2/3 | |
| 2/3 | |
| 10/3 | |
| 10/3 |
Table 2: Average communication times for each dependency
| Task | |
|---|---|
| 22/3 | |
| 6 | |
| 16/3 | |
| 2/3 |
Table 3: Upward rank of each task
Figure 3 shows three valid schedules for the same problem instance. Figure 3a shows the first schedule we validated in the previous section with makespan . Figure 3b shows the schedule that the HEFT algorithm produces with a slightly better makespan of . Finally, Figure 3c shows the best schedule for this problem instance, which has a makespan of just . This is almost half the makespan of the schedule that HEFT (one of the most widely used scheduling algorithms) produces!
Figure 3a: Initial schedule (makespan = 7)
Figure 3b: HEFT schedule (makespan = 6)
Figure 3c: Optimal schedule (makespan = 3.5)
Questions to Consider
- Upward rank has the important property that a task's upward rank is always greater than the upward rank of its dependent tasks. Why is this important?
- What is the runtime of HEFT in terms of , , , and ?
- Why does HEFT perform poorly on the problem instance in Figure 2? Can you think of an algorithm that would do better?
My Research Interests
Task scheduling is a fundamental problem in computer science that pops up everywhere. In this lecture, we formalized the task scheduling problem for heterogeneous task graphs and compute networks with the objective of minimizing makespan (total execution time) under the related machines model. Many other interesting variants of the task scheduling problem exist (see[graham-survey]).
We also learned HEFT, one of the most popular task scheduling heuristic algorithms, and saw a problem instance on which it performs rather poorly. Hundreds of heuristic algorithms have been proposed in the literature over the past decades ([eleven] has nice descriptions of eleven scheduling algorithms). Due to their reliance on heuristics (since the problem is NP-Hard), all of these algorithms have problem instances on which they perform very poorly.
The performance boundaries between heuristic algorithms are not well-understood, however. This is an area of my research. We look at methodologies for comparing task scheduling algorithms to better understand the conditions under which they perform well and poorly.
Figures 4 and 5 depict results from our efforts in this area. Figure 4 shows benchmarking results for 15 scheduling algorithms on 16 datasets. The color represents the maximum makespan ratio (MMR) of an algorithm on a problem instance in a given dataset. The MMR of an algorithm is essentially how many times worse the algorithm performs on a particular problem instance compared to the other scheduling algorithms. For example, on some problem instances in the cycles dataset, the BIL algorithm performs more than five times worse than another one of the 15 algorithms! On other problem instances in the same dataset, however, the algorithm performs well (MMR=1).
Figure 5 shows results from our own comparison method that pits algorithms against each other and tries to find a problem instance where one algorithm maximally underperforms compared to another. Our hope is that by identifying these kinds of problem instances, we can better understand the conditions under which algorithms perform well/poorly.
Figure 4: Benchmarking results for 15 scheduling algorithms on 16 datasets
Figure 5: Adversarial analysis results for 15 scheduling algorithms
Reading List
Our Work
- Jared Ray Coleman and Bhaskar Krishnamachari. 2024. "Comparing Task Graph Scheduling Algorithms: An Adversarial Approach." arXiv:2403.07120. doi:10.48550/ARXIV.2403.07120
- Jared Ray Coleman, Ravi Vivek Agrawal, Ebrahim Hirani, and Bhaskar Krishnamachari. 2024. "Parameterized Task Graph Scheduling Algorithm for Comparing Algorithmic Components." arXiv:2403.07112. doi:10.48550/ARXIV.2403.07112
- Jared Coleman, Mehrdad Kiamari, Lillian Clark, Daniel D'Souza, and Bhaskar Krishnamachari. 2022. "Graph Convolutional Network-based Scheduler for Distributing Computation in the Internet of Robotic Things." In MILCOM 2022, 1070-1075. doi:10.1109/MILCOM55135.2022.10017673
- Tzanis Anevlavis et al. 2022. "Network synthesis for tactical environments: scenario, challenges, and opportunities." Artificial Intelligence and Machine Learning for Multi-Domain Operations Applications IV 12113: 199-206. doi:10.1117/12.2619048
- Jared Coleman, Eugenio Grippo, Bhaskar Krishnamachari, and Gunjan Verma. 2022. "Multi-objective network synthesis for dispersed computing in tactical environments." Signal Processing, Sensor/Information Fusion, and Target Recognition XXXI 12122: 132-137. doi:10.1117/12.2616187
- Daniel D'Souza, Mehrdad Kiamari, Lillian Clark, Jared Coleman, and Bhaskar Krishnamachari. 2022. "Graph Convolutional Network-based Scheduler for Distributing Computation in the Internet of Robotic Things." The 2nd Student Design Competition on Networked Computing on the Edge. https://github.com/ANRGUSC/gcnschedule-turtlenet
- Jared Coleman and Bhaskar Krishnamachari. 2023. "Scheduling Algorithms Gathered: A Framework for Implementing, Evaluating, and Comparing Task Graph Scheduling Algorithms." Technical Report, University of Southern California.
- Jared Coleman. 2023. "SAGA: Scheduling Algorithms Gathered." GitHub repository. https://github.com/ANRGUSC/saga
Theory
- R. L. Graham. 1969. "Bounds on Multiprocessing Timing Anomalies." SIAM Journal on Applied Mathematics 17(2): 416-429. doi:10.1137/0117039
- Jing-Jang Hwang, Yuan-Chieh Chow, Frank D. Anger, and Chung-Yee Lee. 1989. "Scheduling Precedence Graphs in Systems with Interprocessor Communication Times." SIAM Journal on Computing 18(2): 244-257. doi:10.1137/0218016
- Oliver Sinnen. 2007. Task Scheduling for Parallel Systems. Wiley series on parallel and distributed computing. ISBN: 978-0-471-73576-2
- Abbas Bazzi and Ashkan Norouzi-Fard. 2015. "Towards Tight Lower Bounds for Scheduling Problems." In Algorithms - ESA 2015, 118-129. doi:10.1007/978-3-662-48350-3_11
- R.L. Graham, E.L. Lawler, J.K. Lenstra, and A.H.G. Rinnooy Kan. 1979. "Optimization and Approximation in Deterministic Sequencing and Scheduling: a Survey." In Discrete Optimization II, Annals of Discrete Mathematics 5: 287-326. doi:10.1016/S0167-5060(08)70356-X
Scheduling Algorithms
- Haluk Topcuoglu, Salim Hariri, and Min-You Wu. 1999. "Task Scheduling Algorithms for Heterogeneous Processors." In 8th Heterogeneous Computing Workshop, 3-14. doi:10.1109/HCW.1999.765092
- James Blythe et al. 2005. "Task scheduling strategies for workflow-based applications in grids." In 5th International Symposium on Cluster Computing and the Grid, 759-767. doi:10.1109/CCGRID.2005.1558639
- Tracy D. Braun et al. 2001. "A Comparison of Eleven Static Heuristics for Mapping a Class of Independent Tasks onto Heterogeneous Distributed Computing Systems." Journal of Parallel and Distributed Computing 61(6): 810-837. doi:10.1006/jpdc.2000.1714
- Jing-Jang Hwang, Yuan-Chieh Chow, Frank D. Anger, and Chung-Yee Lee. 1989. "Scheduling Precedence Graphs in Systems with Interprocessor Communication Times." SIAM Journal on Computing 18(2): 244-257. doi:10.1137/0218016
- Hyunok Oh and Soonhoi Ha. 1996. "A Static Scheduling Heuristic for Heterogeneous Processors." In Euro-Par '96 Parallel Processing, 573-577. doi:10.1007/BFb0024750
- Andrei Radulescu and Arjan J. C. van Gemund. 2000. "Fast and Effective Task Scheduling in Heterogeneous Systems." In 9th Heterogeneous Computing Workshop, 229-238. doi:10.1109/HCW.2000.843747
- Gilbert C. Sih and Edward A. Lee. 1993. "A Compile-Time Scheduling Heuristic for Interconnection-Constrained Heterogeneous Processor Architectures." IEEE Transactions on Parallel and Distributed Systems 4(2): 175-187. doi:10.1109/71.207593
- A. Poylisher et al. 2021. "Tactical Jupiter: Dynamic Scheduling of Dispersed Computations in Tactical MANETs." In MILCOM 2021, 102-107. doi:10.1109/MILCOM52596.2021.9652937
- Diyi Hu and Bhaskar Krishnamachari. 2019. "Throughput Optimized Scheduler for Dispersed Computing Systems." In 7th IEEE International Conference on Mobile Cloud Computing, 76-84. doi:10.1109/MobileCloud.2019.00018
- R. Armstrong, D. Hensgen, and T. Kidd. 1998. "The relative performance of various mapping algorithms is independent of sizable variances in run-time predictions." In Proceedings Seventh Heterogeneous Computing Workshop, 79-87. doi:10.1109/HCW.1998.666547
- Hesham El-Rewini and T. G. Lewis. 1990. "Scheduling parallel program tasks onto arbitrary target machines." Journal of Parallel and Distributed Computing 9(2): 138-153. doi:10.1016/0743-7315(90)90042-N
- Chung-Yee Lee, Jing-Jang Hwang, Yuan-Chieh Chow, and Frank D. Anger. 1988. "Multiprocessor scheduling with interprocessor communication delays." Operations Research Letters 7(3): 141-147. doi:10.1016/0167-6377(88)90080-6
- Tchimou N'Takpé and Frédéric Suter. 2006. "Critical Path and Area Based Scheduling of Parallel Task Graphs on Heterogeneous Platforms." In 12th International Conference on Parallel and Distributed Systems, 3-10. doi:10.1109/ICPADS.2006.32
Surveys and Algorithm Comparison Papers
- Khushboo Singh, Mahfooz Alam, and Sushil Kumar Sharma. 2015. "A survey of static scheduling algorithm for distributed computing system." International Journal of Computer Applications 129(2): 25-30. doi:10.5120/ijca2015906828
- Essam H. Houssein et al. 2021. "Task Scheduling in Cloud Computing based on Meta-heuristics: Review, Taxonomy, Open Challenges, and Future Trends." Swarm and Evolutionary Computation 62: 100841. doi:10.1016/j.swevo.2021.100841
- T.D. Braun et al. 1999. "A comparison study of static mapping heuristics for a class of meta-tasks on heterogeneous computing systems." In Proceedings Eighth Heterogeneous Computing Workshop, 15-29. doi:10.1109/HCW.1999.765093
- Louis-Claude Canon, Emmanuel Jeannot, Rizos Sakellariou, and Wei Zheng. 2008. "Comparative Evaluation Of The Robustness Of DAG Scheduling Heuristics." In Grid Computing - Achievements and Prospects, 73-84. doi:10.1007/978-0-387-09457-1_7
- Y.-K. Kwok and I. Ahmad. 1998. "Benchmarking the task graph scheduling algorithms." In Proceedings of the First Merged International Parallel Processing Symposium, 531-537. doi:10.1109/IPPS.1998.669967
- Huijun Wang and Oliver Sinnen. 2018. "List-Scheduling versus Cluster-Scheduling." IEEE Transactions on Parallel and Distributed Systems 29(8): 1736-1749. doi:10.1109/TPDS.2018.2808959
- Pooria Namyar et al. 2022. "Minding the gap between Fast Heuristics and their Optimal Counterparts." In Hot Topics in Networking.
- Ashish Kumar Maurya and Anil Kumar Tripathi. 2018. "On benchmarking task scheduling algorithms for heterogeneous computing systems." Journal of Supercomputing 74(7): 3039-3070. doi:10.1007/S11227-018-2355-0
- Jakub Beránek, Stanislav Böhm, and Vojtech Cima. 2022. "Analysis of workflow schedulers in simulated distributed environments." Journal of Supercomputing 78(13): 15154-15180. doi:10.1007/s11227-022-04438-y
- Mohammad Reza Alizadeh et al. 2020. "Task scheduling approaches in fog computing: A systematic review." International Journal of Communication Systems 33(16): e4583.
Machine Learning Approaches
- Mehrdad Kiamari and Bhaskar Krishnamachari. 2022. "GCNScheduler: scheduling distributed computing applications using graph convolutional networks." In Proceedings of the 1st International Workshop on Graph Neural Networking, 13-17. doi:10.1145/3565473.3569185
- Habib Izadkhah. 2019. "Learning based genetic algorithm for task graph scheduling." Applied Computational Intelligence and Soft Computing 2019. Hindawi.
- Penghao Sun et al. 2021. "Deepweave: Accelerating job completion time with deep reinforcement learning-based coflow scheduling." In Proceedings of the Twenty-Ninth International Conference on International Joint Conferences on Artificial Intelligence, 3314-3320.
- Hongzi Mao et al. 2019. "Learning scheduling algorithms for data processing clusters." In Proceedings of the ACM special interest group on data communication, 270-288.
Data and Other References
- M. Rynge et al. 2014. "Producing an Infrared Multiwavelength Galactic Plane Atlas Using Montage, Pegasus, and Amazon Web Services." In Astronomical Data Analysis Software and Systems XXIII, 211.
- Zonghan Wu et al. 2021. "A Comprehensive Survey on Graph Neural Networks." IEEE Transactions on Neural Networks and Learning Systems 32(1): 4-24. doi:10.1109/TNNLS.2020.2978386
- Thomas N. Kipf and Max Welling. 2016. "Semi-Supervised Classification with Graph Convolutional Networks." arXiv:1609.02907.
- Patrycja Krawczuk et al. 2021. "A Performance Characterization of Scientific Machine Learning Workflows." In 2021 IEEE Workshop on Workflows in Support of Large-Scale Science, 58-65. doi:10.1109/WORKS54523.2021.00013
- Daniel Cordeiro et al. 2010. "Random graph generation for scheduling simulations." In 3rd International ICST Conference on Simulation Tools and Techniques. doi:10.4108/ICST.SIMUTOOLS2010.8667
- Takao Tobita and Hironori Kasahara. 2002. "A standard task graph set for fair evaluation of multiprocessor scheduling algorithms." Journal of Scheduling 5(5): 379-394. doi:10.1002/jos.116
- Bin Xiang, Jocelyne Elias, Fabio Martignon, and Elisabetta Di Nitto. 2021. "A dataset for mobile edge computing network topologies." Data in Brief 39: 107557. doi:10.1016/j.dib.2021.107557
- Ewa Deelman et al. 2015. "Pegasus, a workflow management system for science automation." Future Generation Computer Systems 46: 17-35. doi:10.1016/j.future.2014.10.008
- Michael Albrecht, Patrick Donnelly, Peter Bui, and Douglas Thain. 2012. "Makeflow: A Portable Abstraction for Data Intensive Computing on Clusters, Clouds, and Grids." In Proceedings of the 1st ACM SIGMOD Workshop on Scalable Workflow Execution Engines and Technologies. doi:10.1145/2443416.2443417
- Paolo Di Tommaso et al. 2017. "Nextflow enables reproducible computational workflows." Nature Biotechnology 35(4): 316-319. doi:10.1038/nbt.3820
- Scott Kirkpatrick, C. Daniel Gelatt Jr, and Mario P. Vecchi. 1983. "Optimization by simulated annealing." Science 220(4598): 671-680. doi:10.1126/science.220.4598.671
- Ashish Vaswani et al. 2017. "Attention is All you Need." In Advances in Neural Information Processing Systems 30, 5998-6008.