Full text
2013 101 Javier Celaya Alastrué STaRS: a scalable task routing approach to distributed scheduling Departamento Director/es Informática e Ingeniería de Sistemas Arronategui Arribalzaga, Unai Director/es Tesis Doctoral Autor Repositorio de la Universidad de Zaragoza – Zaguan http://zaguan.unizar.es UNIVERSIDAD DE ZARAGOZA
Departamento Director/es Javier Celaya Alastrué STaRS: A SCALABLE TASK ROUTING APPROACH TO DISTRIBUTED SCHEDULING Director/es Informática e Ingeniería de Sistemas Arronategui Arribalzaga, Unai Tesis Doctoral Autor 2013 Repositorio de la Universidad de Zaragoza – Zaguan http://zaguan.unizar.es UNIVERSIDAD DE ZARAGOZA
Departamento Director/es Director/es Tesis Doctoral Autor Repositorio de la Universidad de Zaragoza – Zaguan http://zaguan.unizar.es UNIVERSIDAD DE ZARAGOZA
STaRS: A Scalable Task Routing Approach to Distributed Scheduling PHD DISSERTATION Javier Celaya Alastrué Advised by Unai Arronategui Arribalzaga Departamento de Informática e Ingeniería de Sistemas Escuela de Ingeniería y Arquitectura Universidad de Zaragoza Submitted in fulfillment of the requirements for the degree of Doctor of Philosophy in Computer Science and Systems Engineering with “International Doctor” Mention from Universidad de Zaragoza. June 2013.
iii STaRS: A Scalable Task Routing Approach to Distributed Scheduling Javier Celaya Alastrué Advisor: Unai Arronategui Arribalzaga Universidad de Zaragoza, Spain Reviewers: Arnaud Legrand Université de Grenoble, France Loris Marchal Ecole Normale Supérieure de Lyon, France Committee: José Manuel Colom Piazuelo Universidad de Zaragoza, Spain José Ángel Bañares Bañares Universidad de Zaragoza, Spain Frédéric Vivien Ecole Normale Supérieure de Lyon, France Rizos Sakellariou University of Manchester, United Kingdom Sergio Arévalo Viñuales Universidad Politécnica de Madrid, Spain José Miguel Alonso Euskal Herriko Unibertsitatea, Spain Leandro Navarro Moldes Universitat Politècnica de Catalunya, Spain
Quiero dedicar esta tesis a Cristina y Kira, por su infinita paciencia, amor y cariño. v
Contents List of Figures xv List of Tables xvii List of Algorithms xix 1. Introduction 1 1.1. Motivation . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 1 1.2. Hypothesis and Goals . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 2 1.3. Context and Contributions . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 3 1.4. Thesis Overview . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 4 2. State of the Art 7 2.1. Centralized Designs . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 7 2.2. Decentralization with Aggregated Information . . . . . . . . . . . . . . . . . . . . . . . . . 8 2.2.1. Aggregating Information on Trees . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 9 2.3. Hierarchical Overlays . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 9 3. A Common Scheduling Model for Many Policies 11 3.1. Architecture Overview . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 12 3.1.1. Scheduling Model . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 12 3.1.2. Node Roles . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 13 3.1.3. Overlay Structure . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 15 3.1.4. Fault Tolerance . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 16 3.2. Availability Information Management . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 17 3.2.1. Describing Node Availability . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 17 3.2.2. Availability Aggregation Scheme. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 18 3.2.3. Distributing the Availability Information . . . . . . . . . . . . . . . . . . . . . . . . 23 3.3. Task Routing . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 24 3.3.1. Forwarding Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 24 3.3.2. Routing Patterns . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 25 4. MMP: Makespan Minimization Policy 27 4.1. Makespan Scheduling . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 27 4.2. Local Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28 4.2.1. Measuring the Execution Time . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 28 xiii
xiv Contents 4.3. Global Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 29 4.3.1. Availability Information . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 29 4.3.2. Forwarding Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 30 5. DP: A Policy for Applications with Deadlines 33 5.1. Applications with Time Requirements . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 33 5.2. Local Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34 5.2.1. Availability Function . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34 5.3. Global Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 36 5.3.1. Availability Information . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 36 5.3.2. Forwarding Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 39 6. WDP: DP for Workflow Applications 41 6.1. Scheduling of Workflow Applications . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 41 6.1.1. Workflow Management . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 41 6.2. Local Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 43 6.3. Global Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 44 6.3.1. Availability Information . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 44 6.3.2. Forwarding Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 45 7. FSP: Fair Share Policy 47 7.1. Fairness as the Scheduling Objective . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 47 7.1.1. Slowness . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 48 7.2. Local Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 48 7.2.1. Minimizing the Local Maximum Slowness . . . . . . . . . . . . . . . . . . . . . . . 48 7.2.2. Availability Function . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 50 7.3. Global Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 55 7.3.1. Availability Information Management . . . . . . . . . . . . . . . . . . . . . . . . . . 55 7.3.2. Forwarding Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 57 8. Experimentation 59 8.1. Aggregation Tests . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 59 8.2. Simulation Setup . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 62 8.3. Scalability Results . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 63 8.4. Policy Performance Results . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 66 8.4.1. IBP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 67 8.4.2. MMP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 68 8.4.3. DP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 68 8.4.4. FSP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 70 8.5. Fault-tolerance Results . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 71 8.6. WDP Policy Tests . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 73 8.6.1. Simulation Setup . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 73 8.6.2. Results . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 74 8.7. Comparison with Other Works . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 75
Contents xv 8.8. Discussion . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 76 9. Conclusions 79 Bibliography 81 A. Notation 89 B. The STaRS Simulator 91 B.1. History . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 91 B.2. Design . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 91 B.3. Future Work . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 92 C. Centralized Version of each Policy 93 C.1. IBP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 93 C.2. MMP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 93 C.3. DP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 93 C.4. FSP Policy . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 95 Glossary 97
List of Figures 1.1. Many applications require high amounts of computational resources. For instance, chemical combination models, outer space signal processing and climate change simulations. Pictures shared with a Creative Commons Attribution License by users wasoxygen, karenandbrademerson and NASA Goddard Photo and Video, from Flickr (http://www.flickr.com). . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 2 3.1. STaRS architecture. The scheduling model is divided into a local and a global part. Execution nodes provide the local scheduling, while submission and routing nodes provide the global one. Execution and submission nodes are placed at the leaves of a tree-based overlay, while routing nodes occupy the branches. . . . . . . . . . . . . . . . . 12 3.2. Interactions between the node roles. First, availability information from execution nodes is aggregated by routing nodes. Then, submission nodes issue a request, that is routed to the execution nodes. When tasks are allocated to execution nodes, they communicate directly with the source submission node. . . . . . . . . . . . . . . . . . . . 14 3.3. Mapping of the different node roles. Every physical node plays as a leaf and a branch node in the tree. On top of them, execution and submission nodes are mapped to the leaves, and routing nodes are mapped to the branches. . . . . . . . . . . . . . . . . . . . . . 15 3.4. Hierarchical agglomerative clustering of sampled functions. At each step, the two most similar functions are joined together. . . . . . . . . . . . . . . . . . . . . . . . . . . . . 19 3.5. Path followed by a request that is forwarded towards the execution nodes where it is allocated. At each routing node, it can be divided so that independent branches are explored concurrently. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 26 4.1. The minimum makespan of a random schedule in (a) with two configurations: with (b) divisible tasks and (c) atomic tasks. The former is always shorter or equal to the later. 27 5.1. Example of lu ( δ )for a queue with three tasks. It can be seen that it is a piecewise linear function. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 36 5.2. h.Linterpolates the minimum of f.Land g.L, with sample at s2removed. . . . . . . . 37 6.1. A typical DAG with 10 tasks and several dependencies between them. . . . . . . . . . . 42 6.2. Task queue with four tasks, and the holes available between them. . . . . . . . . . . . . . 43 7.1. Value zi j , where tasks τi and τj switch positions because their deadline functions cross. 49 7.2. A possible task queue with ntasks. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 52 7.3. Example of a zu ( a )function with four pieces. amin is the shortest length a task may have. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 54 xvii
xviii List of Figures 7.4. Two examples of how to join two pieces (black) into one (red). . . . . . . . . . . . . . . . 56 8.1. Aggregation accuracy of a set of 1024 nodes for an increasing SFmax , for the (a) IBP, (b) MMP, (c) DP and (d) FSP policy parameters. The accuracy represents the fraction of the actual availability of a set of nodes that is represented in the aggregated summary. It is very similar for the scalar parameters. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 60 8.2. Aggregation accuracy with 200 sampled functions per summary for an increasing number of nodes, for the (a) IBP, (b) MMP, (c) DP and (d) FSP policy parameters. The accuracy represents the fraction of the actual availability of a set of nodes that is represented in the aggregated summary. It is very similar for the scalar parameters. . . 61 8.3. Average throughput of each policy by network size. It depends on the workload characteristics, like the distribution of the number of tasks and release time. . . . . . . 63 8.4. Average allocation time against network size and policy, of a 1000 task request, with the (a) slow and (b) fast network link, by the decentralized versions of each policy. . . 64 8.5. Average allocation time against network size and policy, of a 1000 task request, with the (a) slow and (b) fast network link, by the centralized versions of each policy. . . . . 64 8.6. Maximum percentage of link bandwidth used by policy and network size, with the fast link model, a sampling interval of 1 second and an SFmax of 200. . . . . . . . . . . . 66 8.7. (a) Finished tasks and (b) finished computation by the IBP policy, on a million nodes, for different values of SFmax , compared to its centralized version and a random allocation. 67 8.8. Performance variation in 5 simulations of the MMP policy for different update bandwidth limit with (a) the slow link model and (b) the fast link model. . . . . . . . . 69 8.9. Makespan by the MMP policy, on a hundred thousand nodes, for different values of SFmax, compared to its centralized version and a random allocation. . . . . . . . . . . . . 69 8.10. (a) Finished tasks and (b) finished computation by the DP policy, on a hundred thousand nodes, for different values of SFmax , compared to its centralized version and a random allocation. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 70 8.11. Maximum slowness among coexisting applications during the simulation, by the FSP policy, on a hundred thousand nodes, for different values of SFmax , compared to its centralized version. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 71 8.12. Finished computation for simulations with churn, with median session times of 60, 30, 15 and 5 minutes, for the (a) IBP, (b) MMP, (c) DP and (d) FSP policies. It is normalized to the results without churn. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 72 8.13. A Fork-Join graph model with 10 tasks. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 73 8.14. A Laplace equation solver graph model with 9 tasks. . . . . . . . . . . . . . . . . . . . . . . 73 8.15. Allocation time for different workflow widths, in a network of 1 million nodes. . . . . 74
List of Tables 6.1. Interval endpoints generated from three example creation times t0 . The difference between each endpoint and t0 is always between one and two times the intended interval duration. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 45 8.1. 99th percentile and maximum percentage of link bandwidth used by link model, policy, sampling interval and SFmax , on simulations with an update bandwidth limit of 100 KBps. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 65 8.2. Average speedup by workflow model and priority. . . . . . . . . . . . . . . . . . . . . . . . 75 8.3. Comparison of several distributed scheduling projects. . . . . . . . . . . . . . . . . . . . . 76 A.1. Notation in text and algorithms. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 89 xix
List of Algorithms 3.1. Aggregation of two availability summaries. . . . . . . . . . . . . . . . . . . . . . . . . . . . . 21 3.2. Forwarding algorithm for the IBP policy. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 25 4.1. Forwarding algorithm for the MMP policy. . . . . . . . . . . . . . . . . . . . . . . . . . . . . 31 5.1. Computation of lu(δ)on node Pu.. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 35 6.1. Extract the sequences of dependent tasks from G. . . . . . . . . . . . . . . . . . . . . . . . 43 6.2. Generate a set of ninterval endpoints from the current time ν.. . . . . . . . . . . . . . . 44 6.3. Forwarding algorithm for the WDP policy. . . . . . . . . . . . . . . . . . . . . . . . . . . . . 46 7.1. Find the set of boundary values for the tasks in the queue Q.. . . . . . . . . . . . . . . . 50 7.2. Sort a task queue to minimize its maximum slowness. . . . . . . . . . . . . . . . . . . . . . 51 7.3. Recalculate deadlines for a maximum slowness and sort tasks. . . . . . . . . . . . . . . . . 51 7.4. Check if all tasks in a queue meet their deadlines for a certain slowness. . . . . . . . . . 51 7.5. Computation of zu(a)pieces. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 53 7.6. Forwarding algorithm for the FSP policy. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 58 7.7. Get the minimum slowness that can be reached when assigning the tasks in request to the nodes described by info.. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 58 C.1. Centralized version of the IBP policy forwarding algorithm. . . . . . . . . . . . . . . . . . 94 C.2. Centralized version of the MMP policy forwarding algorithm. . . . . . . . . . . . . . . . 94 C.3. Centralized version of the DP policy forwarding algorithm. . . . . . . . . . . . . . . . . . 95 C.4. Centralized version of the FSP policy forwarding algorithm. . . . . . . . . . . . . . . . . 96 xxi
Chapter 2. State of the Art STaRS provides a scalable scheduling model to virtually any type of distributed computing platform, from loosely coupled desktop grids to dense clusters. Its design decisions are motivated by the shortcomings we have observed in the related literature. Centralized designs have a structurally limited scalability and resilience, so many decentralized alternatives have been proposed. Most use aggregated information in some way, but they fail to provide good scalability, fault-tolerance and flexibility at the same time. 2.1. Centralized Designs Centralized scheduling services are common in cluster, grid and cloud computing, so their limitations are well known. The most popular desktop grid platform is BOINC [ 8 ], with millions of volunteers in projects like SETI@home [ 9 ]. However, it reaches these high scales by having the global scheduler manage very little detail about each execution node. On the contrary, Condor [ 90 ]provides a powerful scheduling engine to cluster and grid platforms at the expense of needing ad-hoc solutions [ 43 ] to reach scales over the thousand nodes. Falkon [ 77 ]aims for higher scales sacrificing some of the functionality provided by Condor, like multiple queues and priorities. However, its centralized dispatcher still scales to a maximum number of managed tasks and resources. In their tests, they reach two million tasks on 54,000 nodes. Something similar happens with Mesos [ 50 ]. It does not perform any scheduling, it just manages a set of resources and provides a fair share among several distributed computing frameworks. It offers resources to the platforms, and they perform the actual scheduling. In this way, it should be able to scale to more than the 50,000 nodes the authors advertise, but a maximum is inevitable. STaRS circumvents these limitations because it aggregates rich information about execution nodes to avoid centralizing the scheduling decisions. It reaches higher scales than Condor or Falkon, without an expected maximum, while being able to implement a wider range of scheduling policies than BOINC. Other platforms with application-specific policies rely on centralized resource management. Nimrod [ 6 , 4 , 5 , 3 ]performs parameter exploration experiments acting as a resource broker: it first gathers the available resources that fit the experiment requirements and then schedules the tasks on them. MapReduce [ 39 ]and Hadoop [ 73 ]process large-scale data sets, centralizing data management and job scheduling on dedicated nodes. Policies based on economic models [ 19 ]set prices and offers independently from each other, but inflation adjustment and price gathering problems have been solved in practice with centralized market managers [ 48 , 21 ]or non-scalable algorithms [ 95 , 93 ]. Many-task computing platforms [ 76 , 53 ]focus on the scheduling of a large number of short, data-intensive tasks. Raicu et al. [ 76 ]think that they must relax some constraints found in other platforms to be able to 7
8 State of the Art cope with this amount of fine-grained workload. This is the case of Falkon commented before. All these works would benefit from a decentralized model like STaRS. 2.2. Decentralization with Aggregated Information Several proposals adapt grid schedulers to use aggregated information describing each domain. They are still not fully decentralized, but claim to increase their scalability by using less detailed information with similar capabilities. Rodero et al. [ 84 ]use this information to rank brokers of each domain by suitability. They define a distributed meta-broker architecture that distributes the aggregated information to every domain, so that the best broker for a job can be selected. Kokkinos and Varvarigos [ 59 ]use aggregation to reduce the amount of information exchanged between domains, limit the exposure of sensitive information and improve interoperability. They focus on the aggregation accuracy for different kinds of attributes and clustering operators. Unfortunately, the architecture is only a two-level hierarchy, with a centralized meta-scheduler at the top level, that decides which domain a task is allocated to. Brunner et al. [ 18 ]propose the use of a conceptual clustering algorithm to summarize resource capabilities in different grid domains. They focus on finding suitable execution nodes in near locations to reduce data transmission times. There is no centralized scheduler, but all domains must know each other to maximize the matchmaking performance. Rahman et al. [ 75 ]map a d-dimensional logical index to a distributed hash table (DHT). Execution nodes are published in the index, using capabilities as coordinates, and they are found with point and range queries. Many decentralized solutions have been proposed to date, some of them also using aggregated information. In the field of the resource discovery with aggregated information, Cai and Hwang [ 23 ] build distributed aggregation trees on a Chord-like DHT[ 88 ], that aggregate node information for Grid resource monitoring. SWORD [ 7 ]allows complex queries to search for computational resources. The authors propose a decentralized implementations based on partitioning the availability information space into ranges and map them to a DHT. Each node is responsible for the availability information of the set of resources in the same range. However, partitioning must be done carefully to reach good load-balancing. Cohesion [ 85 ]is a decentralized, tree-based information aggregation system on an unstructured peer-to-peer grid platform. They define two ways of building the tree from the set of nodes, one focused on efficiency and another focused on scalability. NodeWiz [ 12 ] implements a Grid information service specialized on range queries. It uses a binary balanced tree to partition an attribute space and is able to look for nodes that match certain criteria, but it is limited to scalar attributes. Cardosa and Chandra [ 26 ]present an hierarchical aggregation method that clusterizes resource capacity distribution functions into so called “resource bundles”. These bundles provide statistical information about resource properties, with a level of confidence. This implies that there is a certain probability of discovering nodes that may not fulfill the requirements of the requested application. Other works also try to decentralize the scheduling component, but we have seen no other work that combines scalability, fault tolerance and flexibility as we do. Diet [ 28 , 27 ]uses a static, ad-hoc hierarchy to route requests to execution nodes, but no aggregated information is used for forwarding. Instead, the root node sends each request to all the execution nodes, and they respond whether they can execute it. Intermediate nodes use stateful queues to control the flow of requests. This architecture is difficult to scale to millions of nodes. It is not very resilient, either, but they maintain
State of the Art 9 links to some ancestors other than the father to manage failures. WaveGrid [ 98 ]uses a CAN [ 80 ] overlay to organize nodes by timezone. It allocates work to idle nodes whose timezone is currently in the night. Resources are discovered by randomly selecting some initial nodes and then using an expanding ring search. While it scales fairly well, the information used for scheduling is too limited and random to allow more complex policies. Kim et al. [ 58 ]also use CAN to organize nodes by their resource availability and aggregate queue length information. They justify this overlay by its superior scalability and fault-tolerance properties over a tree overlay, but its non-hierarchical structure is less suited for the aggregation of data. Their solution is, at each cell of the CAN overlay, to aggregate the information coming only from the same row in each dimension, which discards most information to make scheduling decisions. Kwan and Mupala [ 61 ]use a super-peer network to schedule bag-of-tasks application in volunteer desktop grids. A gossip protocol scatters aggregated information about computing power and location of resources between super-peers. The scheduling policy consists in just assigning tasks to the fastest idle node in the selected super-peer. Unstructured networks are more resilient to failures, as long as super-peers have a very low failure rate. 2.2.1. Aggregating Information on Trees All these works highlight that aggregating resource information in a tree is a popular approach to decentralize the resource discovery process. It scales well by reducing the amount of information that needs to be managed. A similar approach is often found in sensor networks [ 52 , 71 ], where less transmitted information between nodes leads to less consumed energy. Several general-purpose aggregation and indexing frameworks also use hierarchical aggregation. Astrolabe [ 82 ]uses a gossip protocol inside a user-defined hierarchy, with arbitrary aggregation functions. In [ 38 , 74 ], the authors describe a generic aggregation protocol for network management purposes. It creates a tree structure on top of an overlay network, and aggregates network state variables – bandwidth, delay, number of nodes – using operators like sum, average and min/max. Mortar [ 65 ]focuses on data stream management. For each query, it builds several static tree overlays to provide resilience and faster aggregation. It is also very common to build a virtual tree on top of a more resilient structure, like a DHT. That is the case of SDIMS [ 96 ], which improves the tree performance by treating readdominated and write-dominated parameters separately. These works inspired our aggregation scheme, but there are important differences. We are interested in the individual state of each execution node. For this reason, the aggregation is not actually performed until enough samples are collected. Then, a clustering algorithm selects which samples should be aggregated together. This scheme prioritizes the accuracy of the aggregation in the lowest levels, where the forwarding algorithm needs to be more precise. 2.3. Hierarchical Overlays The use of a hierarchical overlay network is important in our design. Both the aggregation scheme and the forwarding algorithm assume that nodes are organized in a binary tree (see Section 3.1.3). We also expect the overlay to recover from node failures. Many works propose fault-tolerant hierarchical overlays. NodeWiz builds a k-d-tree in which nodes are sorted by their properties. Each branch divides the search space in two parts by one of the properties. Its objective is to maintain the tree and the load of its nodes well-balanced. It also considers several failure situations and their solution.
10 State of the Art P-Grid [ 2 ]builds a virtual trie structure on top of a DHT, for storage and search of data. By using self-organization, the tree balances its load even with non-uniform key distributions. However, it does not take fault-tolerance into account. VBI-Tree [ 54 ]builds a virtual balanced tree on top of a DHT. It uses multi-dimensional indexing that allows range queries to be performed with cost O ( log2N ). There are references to failure management, but its mechanisms are not properly described. Finally, Caron et al. [29]contribute a distributed lexical placement table (DLPT), which is a trie that allows exact, partial and range queries. In this case, it uses replication to provide resilience.
Chapter 3. A Common Scheduling Model for Many Policies “All models are false but some models are useful.” — George E. P. Box In this chapter we present the details of STaRS. It is an online distributed scheduling model, customizable with different policies: •Distributed scheduling: STaRS receives request from users for the execution of applications. Each application consists of a set of tasks, and STaRS allocates them to many, independent execution nodes. •Online: Applications are allocated as they arrive to the system, they are not known in advance. •Different policies: A scheduling policy defines the objective of the allocation of tasks: fulfill certain requirements, finish them as soon as possible, provide fairness to users, etc. It is possible to implement different policies on STaRS with a common architecture. •Model: STaRS is a model that describes a set of tools, protocols and algorithms. When implemented on a distributed scheduling platform, it provides with scalability, fault-tolerance and versatility. STaRS tries to fill a gap. It simultaneously provides with the scalability, fault-tolerance and versatility properties that most distributed computing platforms lack. Scalability, as its ability to manage as many execution nodes, users and applications as needed. Fault-tolerance, as its ability to gracefully degrade its performance when part of its components fail, and later recover. Versatility, as its ability to schedule different application types, with different objectives and in different environments. We have seen no other distributed computing platform that provides all of them at the same time. To achieve these goals, we build up STaRS on the principles of decentralization and partial knowledge. In this way, we avoid the bottleneck and single point of failure of a centralized design, and no algorithm needs to know the full state of the system at once. Then, we propose three generic tools: a data model to represent the availability of execution nodes; an aggregation scheme that propagates the availability information on a tree-based overlay network, and that can be tuned to find a tradeoff between its accuracy and cost; and a forwarding algorithm that, using that information, routes tasks towards the most suitable execution nodes, performing stateless resource discovery and allocation 11
12 A Common Scheduling Model for Many Policies FAULT TOLERANCE Leaf Branch TREE OVERLAY Execution Node Submission Node Routing Node NODE MODEL Local Scheduling Global Scheduling SCHEDULING MODEL Figure 3.1.: STaRS architecture. The scheduling model is divided into a local and a global part. Execution nodes provide the local scheduling, while submission and routing nodes provide the global one. Execution and submission nodes are placed at the leaves of a tree-based overlay, while routing nodes occupy the branches. simultaneously. Many policies can be implemented by providing a suitable instantiation of these three elements. In this chapter we present an example policy, the Idle/Busy Policy (IBP), that allocates tasks to idle nodes. In order to facilitate the reading of this and the following chapters, Appendix Acontains a brief reference of the most common notation used throughout it. It also presents the dot notation used in algorithms to refer to the fields of compound values. Nevertheless, notation is further detailed as it is used in the thesis. In particular, the notation that only appears in certain chapters. 3.1. Architecture Overview Figure 3.1 presents STaRS architecture. From top to bottom, the scheduling model consists of a local and a global part. The local scheduling manages the task execution, while the global scheduling manages the resource discovery and task allocation. The functionality of each part is implemented through the node model. Each physical node may play up to three roles: Execution nodes provide the local scheduling, while submission and routing nodes provide the global scheduling. These nodes are organized in a logical tree-based overlay, so that routing nodes occupy the branches and execution and submission nodes hang from the leaves. Finally, the fault tolerance is taken into account at every level. The tree overlay must recover from physical node failures. Then, the node roles that a failed physical node was playing can be reassigned to another one, and the scheduling model can maintain its functionality, with a proportional degradation if needed. 3.1.1. Scheduling Model STaRS presents a common scheduling architecture for different policies. Among other things, the policy defines the type of application that is being scheduled. We have mainly focused on policies for
A Common Scheduling Model for Many Policies 13 bag-of-tasks applications, but other kinds of application could be accepted by the system. For instance, the workflow of tasks is commonly found in many scientific applications. Chapter 6presents a policy for workflow applications. A bag-of-tasks application Ai consists of a set of ni independent, identical tasks. Additionally, depending on the policy being used, an application may have other parameters that describe the properties and requirements of its tasks. Common ones are the length of the tasks ai , measured in millions of floating point operations (FLOPs), and their required memory and disk space, mi and di , measured in megabytes. Let PRi = ( ni,ai,mi,di,... )be the tuple of parameters that describes application Ai . Then, a user submits an applications to the system as an application scheduling request containing PRi . This kind of configuration is very common in distributed scheduling platforms due to its high degree of parallelism. The scheduling model is divided into a local and global part. The former describes the local scheduler and its policy. The local scheduler of a node manages a queue with the received remote tasks. The local policy decides whether new tasks can be accepted into the queue, and their order of execution. Common examples of local policies are First Come First Served (FCFS) or Earliest Deadline First (EDF). The global part describes the availability information, the aggregation scheme and the forwarding algorithm. They are deeply discussed in Sections 3.2 and 3.3. Every local scheduler periodically calculates its availability to execute different types of applications, and exports this information. The aggregation scheme distributes it among the nodes of the system in a hierarchical fashion. It is then used by the forwarding algorithm to route application scheduling requests towards the most suitable execution nodes. The global scheduling policy determines the implementation of these three elements. Naturally, the global and local policies must match. For instance, a global policy which tries to fulfill application deadlines will be used along with an EDF local policy. Thus, throughout this thesis we refer to both of them as just the scheduling policy, without distinction. Five common policies are presented in the following chapters. 3.1.2. Node Roles Every physical node Pu of the system plays three node roles that provide the scheduling model functionality: • The execution node Eu contains a local scheduler and an execution environment. The execution environment provides the platform-dependent mechanisms for the safe execution of remote tasks. • The submission node Su is the interface between the user and the platform. It manages the submission of application requests, and monitors the activity of any remote task that has been successfully allocated to an execution node. • The routing node Ru is the component that implements the global scheduling policy rules. First, it receives the availability information from the execution nodes and distributes it to its neighbors. Then, it forwards application requests towards the most suitable execution nodes using that information.
14 A Common Scheduling Model for Many Policies R R R SE 2 - Submit 3 - Forward 3 - Forward 4 - Allocate 5 - Task I/O 1 - Aggregate 1 - Aggregate Figure 3.2.: Interactions between the node roles. First, availability information from execution nodes is aggregated by routing nodes. Then, submission nodes issue a request, that is routed to the execution nodes. When tasks are allocated to execution nodes, they communicate directly with the source submission node. These roles are organized in a hierarchical fashion (Figure 3.2). The local scheduler at each execution node sends its availability information to its father routing node. Routing nodes aggregate this information through the tree. When a user wants to get an application executed, it instructs a submission node to issue an application scheduling request. This request is forwarded by the routing nodes towards the most suitable execution nodes, using the aggregated information. Execution nodes are allocated as they are found, in a single stage. From then on, task communication is performed directly between the execution node and the source submission node. Usually, all physical nodes play these three roles in order to balance the load among them, but other configurations may also be interesting, e.g. only some nodes playing the submission node role in a dedicated cluster. The roles played by the same physical node are placed independently within the overlay. This provides great flexibility by allowing the relocation of any of the roles of a node without affecting the other ones. The execution environment at each execution node Eu isolates remote tasks from the rest of the system. A virtual machine is a straightforward implementation, but more lightweight ones can be found in some platforms (e.g. Linux containers). This is done mainly for security reasons, to prevent remote tasks from abusing the host computers. It also allows the node owner to arbitrarily limit the type and amount of resources that remote tasks may use. These values are often used to calculate the availability of an execution node, most common ones being: •The computational power su, measured in millions of FLOPs per second. •The available memory reserved for the execution of remote tasks Mu, measured in megabytes. •The available disk space Du, measured in megabytes too. In order to be able to execute every remote task under the same conditions, we assume that the environment disallows preemption. This feature is commonly found in existing distributed computing platforms, since it prevents preempted tasks from consuming memory and disk space to store their
A Common Scheduling Model for Many Policies 15 PHYSICAL NETWORK TREE OVERLAY NODE MODEL R R R SE Figure 3.3.: Mapping of the different node roles. Every physical node plays as a leaf and a branch node in the tree. On top of them, execution and submission nodes are mapped to the leaves, and routing nodes are mapped to the branches. paused state. So, at each execution node, there is only one running task at a time, which either finishes or is aborted. 3.1.3. Overlay Structure STaRS assumes the existence of an underlyingoverlay network with a balanced binary tree organization. In fact, it assumes a design similar to VBI-Tree [ 54 ]. It is a search tree that differentiates branches from leaves. Leaf nodes contains the resources that the system looks for. Branch nodes route discovery requests to find the matching resources. Physical nodes play both roles, so that they can contain resources and route requests at the same time. In our case, physical nodes play the routing node role at the branches and the execution and submission node roles at the leaves (Figure 3.3). Every physical node may play all three roles in different positions of the tree, mutually independent from each other. The tree is balanced so that the management, search and distribution of data operations are performed with cost O ( log2N ). This provides very good scalability properties to the system. The tree is binary because, from our experience in [ 32 ], it has a better tradeoff between computational cost and tree height than higher-degree trees. On one hand, at each node, the network traffic and the cost of most core algorithms are proportional to the number of children: maintain tree links, receive
16 A Common Scheduling Model for Many Policies and aggregate children information, route tasks to child nodes, etc. On the other hand, increasing the tree degree decreases the tree height, and so it shortens the path between any two nodes. This reduces the time of two important processes: routing a request to the execution nodes and distributing the availability information through the tree. However, for every tree degree over 3, while the computational cost increases by a factor of k , the tree height is reduced by a factor of less than k . So, there is no performance benefit if we increase the tree degree over 3. Then, it is much simpler to manage a binary tree than a 3-degree tree, so the binary option is preferred. We also assume that nodes are ordered in the tree, grouping nearby execution nodes under the same branch. By doing so, certain locality information is provided to the forwarding algorithm in order to look first for execution nodes nearer to the submission node. This can be accomplished with topology-aware overlay construction, like in [ 81 ]. However, it does not meaningfully contribute to the evaluation of STaRS, so we use the simpler method of sorting nodes by network address in our tests. In order to bootstrap the overlay, we propose using a classical approach in decentralized systems, where a set of well-known nodes are used as a persistent entry point. We proposed a first design of such an overlay in [ 32 ]. It provided the methods to build the tree structure and expose it to the scheduling components, but it had limited fault-tolerance capabilities. Later, two master thesis [ 31 , 67 ]have been carried out to overcome these limitations, taking two different approaches. The first one backs up the tree structure on a DHT, and the second uses redundant links between nodes. They show the feasibility of building a scalable and fault-tolerant tree-based overlay, being able to successfully recover from multiple node failures, but they still need further development. For this reason, and since this thesis is focused on the scheduling part, the overlay behavior is simulated in the experiments. 3.1.4. Fault Tolerance The management of faults in STaRS takes a best-effort approach. It tries its best to allocate and finish applications, but failures are admissible and should be expected by users. A failed node can make requests reach no execution node, loose the availability information of its branch and abort the tasks it was executing. Thus, every level of the model must consider resilience to faults. First, the overlay network must be fault-tolerant in order to provide a reliable structure. Several peer-to-peer overlays already exist that construct a tree structure with good scalability and faulttolerance properties [ 12 , 54 , 2 ]. Additionally, we have the experience of the two master thesis commented in the previous section. All of them prevent the top levels of the tree from turning into a bottleneck and a single point of failure. So, in our tests, we assume that the overlay is able to recover from node failures, and we evaluate the impact of the recovery in the scheduling performance. We do not evaluate the cost of the recovery, as the authors of each referred overlay have already done it. For failures in the node and scheduling models, we treat routing, submission and execution nodes independently. Routing node failures affect the aggregation scheme and the forwarding algorithm. A failed routing node looses the availability information aggregated from its branch. The availability information is distributed in a reactive way, it is sent whenever a routing node modifies its information or a link with a neighbor changes. Likewise, when a routing node detects that a neighbor has failed, it sends an update as soon as the link is recovered. Any missing information is quickly rebuilt. A failed routing node also disconnects its branch from the rest of the tree. Scheduling requests cannot jump from one side to the other until the node recovers, so we assume that this will impact
A Common Scheduling Model for Many Policies 23 Now, the clustering algorithm is able to take the two sampled functions that produce the minimum MSE when summed up. With one scalar parameter, the distance operation of two sampled functions is straightforward: it just returns the MSE of their sum. But with more than one scalar parameter, we have many MSE values as the result of the sum operation, one for each parameter. As we explained before, the best solution is to calculate the distance as a weighted sum of the normalized MSE for each parameter. That is, for nscalar parameters p1through pn, DIST(f,g) = n X i=1 αiNORMi(SUM(f,g).msepi). (3.8) The αi coefficients represent the weight of each parameter in the distance operation. Usually, they will all have the same value, but in some situations it is interesting to give more importance to some of the parameters over the rest. For the IBP policy, we have decided that both memory and disk space have the same weight. Then, the normalization function NORMi that appears in (3.8) maps the MSE of parameter pi to the [0 , 1]range. Let B be the set of all the execution nodes of the current branch, whose information is being clustered, then NORMi(m) = m max E∈B(E.pi)−min E∈B(E.pi)2. (3.9) Values maxE∈B ( E.pi )and minE∈B ( E.pi )must be known at each routing node, so they are also aggregated. An availability summary contains a maximum and a minimum value for each parameter pi . Execution nodes initialize them with their own value. When two summaries are aggregated, the resulting summary simply gets the maximum of the maximums and the minimum of the minimums. So, the cost in time and space is very little. If the maximum and minimum values for any parameter are equal, the sum operation cannot make an error on that parameter, and it is not taken into account in the linear combination. 3.2.3. Distributing the Availability Information The aggregation of the availability summaries is done in a hierarchical fashion. Execution nodes send a summary to their father routing node with a single sampled function, representing their own availability. Then, routing nodes aggregate the summaries from both children, and send the result further up the tree. This scheme is reactive: whenever execution nodes change their availability or routing nodes receive new information, the aggregation mechanism is triggered. Execution node availability changes periodically. For some policies, it may change only every time a task starts or finishes. In other cases, it may be frequently changing, since time can be an important factor of the execution node state. For this reason, the availability information must be kept up to date at the routing nodes. However, if many execution nodes change at the same time, a cascading effect will flood the upper levels with update summaries, rendering the system unusable. To provide scalability and reliability, this problem must be faced. We propose two solutions. The first one is implicit to the aggregation scheme. Routing nodes periodically receive updated information from their children nodes, aggregate it and store the result as the information of their branch. But when previous information exists, it is compared with the
24 A Common Scheduling Model for Many Policies new one. If they are equal, it is not reported to the father, since it is the same it already has. This may happen when a child node sends a similar summary to the one it sent before, due to the way in which the sum of sampled functions is performed. If it uses operations like the minimum or maximum, as in the IBP example policy, little variations will probably produce the same result. So, a summary going up the tree may stop before reaching the root when the availability changes lightly. The second measure is to limit the bandwidth used to send summaries. Routing nodes insert a short delay after sending each summary to keep the average used bandwidth under an arbitrary maximum. After the delay expires, only the most recently aggregated summary is sent. This can be done with the availability information because each new summary makes the previous ones obsolete. This bandwidth limit, along with the size of the availability summary, is a tradeoff between the traffic supported by nodes and the time needed by a change in the leaves to reach the root of the tree. The lower the bandwidth limit, the higher the period between two summaries are sent. This results in the availability information being out of date more often. As we will show, this impacts the performance of some policies that are more sensible to availability changes. 3.3. Task Routing Task routing is the process of forwarding the tasks in an application scheduling request towards the most suitable nodes, given its parameters. We call it “task routing” because it resembles the routing of packets in a computer network. A network router uses the forwarding table to decide in which direction to send a packet. Our forwarding algorithm uses the availability information to decide in which direction a request may reach the best execution node for each of its tasks. The task routing process starts when a submission node Su issues a new application scheduling request req for req.n tasks of application Ai . As we explained in Section 3.1.1, we focus on bag-of-tasks applications. For such an application, a request contains the application properties PRi , the address of the requester Su , a request identifier and an interval of task identifiers [1 ,req.n ]. All the tasks in a bag-of-tasks application are identical, so we can group them in an interval instead of enumerating each of them. Then, routing nodes invoke the forwarding algorithm on the requests they receive. 3.3.1. Forwarding Algorithm Using the availability information, the forwarding algorithm decides in which direction it sends each task. So, it ends up sending nl tasks to the left child, nr tasks to the right child and nf tasks to the father, so that nl + nr + nf = req.n . With these values, it creates three new requests with the same application properties, request identifier and requester address as the original. The request for the left child will contain the task identifiers in [1 ,nl ]; the one for the right child, task identifiers in [ nl +1 ,nl + nr ]; and the one for the father, task identifiers in [ nl + nr +1 ,req.n ]. The resulting requests are sent in the corresponding direction, as long as they carry at least one task. However, if there is no father because the current routing node is the root, the nf tasks that were meant for it are discarded. They are treated as a scheduling failure and will be resent by their submission node. The algorithm is stateless with regard to the requests it forwards: it needs to hold no record of them. This improves its scalability and fault tolerance. Each policy must specialize the forwarding algorithm, as it specializes the availability information. Usually, a policy-dependent objective function sorts sampled functions by the suitability of allocating
A Common Scheduling Model for Many Policies 25 Algorithm 3.2 Forwarding algorithm for the IBP policy. Pre: Ruis this routing node. request is the request. Post: reqLeft , reqRight and reqFather are the resulting requests to be sent to the left child, right child and father nodes, respectively. 1: procedure FORWARD(request) 2: availableSF ←; 3: if request.srcAddr 6=Ru.leftAddr then 4: GETSF(Ru.leftInfo, request.PR, availableSF) 5: end if 6: if request.srcAddr 6=Ru.rightAddr then 7: GETSF(Ru.rightInfo, request.PR, availableSF) 8: end if 9: SORT(availableSF).Best nodes are allocated first. 10: while ¬ISEMPTY(availableSF)Vrequest.n>0do 11: sf ←POPFRONT(availableSF) 12: numTasks ←AF(sf, request.PR) 13: if ISFROMLEFTCHILD(sf)then 14: reqLeft ←EXTRACT(request, numTasks) 15: else 16: reqRight ←EXTRACT(request, numTasks) 17: end if 18: end while 19: if request.n>0then 20: reqFather ←request 21: end if 22: end procedure a task to one of their nodes. Then, tasks can be assigned to each child from the most to the least suited sampled function. Besides, in most policies, the forwarding algorithm also updates the availability information of the children branches, trying to estimate how it will change once the sent tasks are allocated. This avoids sending tasks of subsequent requests to the same execution nodes before their availability is updated through the aggregation scheme. For instance, Algorithm 3.2 shows the forwarding algorithm of the IBP policy. Procedure GETSF fills the list availableSF with the sampled functions of those nodes that are able to execute tasks from the request. Then, this list is sorted to allocate first those nodes whose available memory and disk space are closest to the requirements. With this heuristic, nodes with more available resources are saved for future applications with higher requirements. Tasks that cannot be allocated in this branch are sent up to the father routing node. 3.3.2. Routing Patterns Specializing the forwarding algorithm for each policy may produce several patterns. A common one is the IBP policy pattern, which assigns as many tasks as possible to the current branch and send the
26 A Common Scheduling Model for Many Policies ... R R R S... ... R ... R E E n1 n2n3 req.n req.n−n1 req.n−n1−n2req.n−n1−n2 req.n−n1−n2−n3 n4req.n−n1−n2−n3−n4 Figure 3.5.: Path followed by a request that is forwarded towards the execution nodes where it is allocated. At each routing node, it can be divided so that independent branches are explored concurrently. rest to the father. At the next level, the forwarding algorithm assumes that it cannot send any more tasks in the direction the request came from, so it sends some tasks to the other one and the rest again to the father. But if the request comes from the father, all the tasks are sent to the children nodes. Figure 3.5 illustrates this pattern with an example. A request with K tasks is issued by a submission node from a leaf of the tree. At each routing node, the forwarding algorithm separates the ni tasks that can be allocated to its branch, and the rest is sent up. The process continues until all tasks are allocated or all nodes are reached. A similar pattern consists in sending all the tasks to the father node as long as a certain criteria is not met. Then, send them down again, assigning tasks to both children. But all the patterns have some advantages in common. The first one is that the task routing process looks for execution nodes in independent branches concurrently, so its cost depends on the number of network jumps of the longest path. In a balanced binary tree, this cost is O ( log2min ( ni,N )) jumps, where ni is the number of tasks to allocate and N is the size of the network. This cost scales very well with both parameters. Another one is that all the submissions originate in the leafs, and only climb up the tree until they discover enough execution nodes. Doing so, maintains most request traffic in the lower levels of the tree, and uses more accurate availability information. Furthermore, if the execution nodes are ordered in the tree by location, as suggested in Section 3.1.3, the discovered ones will be nearby the requester. This might be helpful when the applications have big input or output data.
Chapter 4. MMP: Makespan Minimization Policy “It’s the job that’s never started takes longest to finish.” — J. R. R. Tolkien 4.1. Makespan Scheduling When the users have no other priority, one of the most common policies consists in trying to finish all the scheduled work in the platform as soon as possible. As shown in Figure 4.1, it is well known that this objective is accomplished when all the nodes finish at the same time. Otherwise, nodes that finish earlier would be able to do part of the work of the nodes that finish later. However, this is only feasible if tasks can be arbitrarily divided. In an heterogeneous environment as we consider, each execution node runs tasks at different speed and with an indivisible amount of computation. So, the objective of such a policy is to minimize the makespan: the maximum time needed by any node to finish its work. This subject has been widely studied since long ago. The problem is usually divided into offline and online scheduling, and considering machines with both identical and different speeds. The optimal offline solution is NP-hard [ 46 ]in any case, but several polynomial time approximation schemes have been proposed. The Graham’s well-known list scheduling algorithm [ 47 ]is (2 − 1 /m )-competitive on m identical machines with linear cost. It sorts the list of jobs by priority and sends each one to the machine that has been assigned the least amount of work so far. Hochbaum and Shmoys [ 51 ] propose another solution with arbitrary relative error. On machines with different speed, Lenstra et al. [ 63 ]give a 2-competitive solution based on integer programming. On the other hand, online (a) (b) (c) Figure 4.1.: The minimum makespan of a random schedule in (a) with two configurations: with (b) divisible tasks and (c) atomic tasks. The former is always shorter or equal to the later. 27
28 MMP: Makespan Minimization Policy solutions deal with the problem of having to decide on each job before knowing about the next one. Fleischer and Wahl present MR [ 45 ], an online algorithm for identical machines, that reaches 1.9201competitiveness, setting the current best upper bound. Englert et al. [ 42 ]propose a 2-competitive online algorithm for machines with different speeds, by buffering the incoming jobs and reordering them. Here, we present the Makespan Minimization Policy (MMP) for STaRS. It is an online, decentralized policy to minimize the makespan among all the currently scheduled applications. In this policy, an execution node queue may hold as many tasks as needed, so the makespan is the maximum end time of the longest queue. Then, similarly to the list scheduling algorithm, tasks are routed to those nodes whose queue will remain the shortest after allocating them. 4.2. Local Policy Unlike the IBP policy, the MMP local schedulers have an unlimited task queue. They accept every incoming task, as long as memory and disk space requirements are met. Then, tasks are inserted at the back of the queue, in FCFS order, because we assume that there is no specific application priority. Since we want to limit the queue length by allocating a similar workload to every node, we introduce a constraint to the queue length in the availability function: it cannot be longer than a certain time. We define the availability function for this policy as AFu ( mi,di,ai,qi ). We recall that mi , di and ai are the required memory, required disk space and length of a task of application Ai , respectively. Then, the function returns the maximum number of tasks of application Ai that can be added to the queue of execution node Eu so that it finishes no later than time qi . This function allows the forwarding algorithm to decide how many tasks can be sent to a node in order to increase its queue to a certain length. Let the queue end time of node Eu be Qu , the availability function is calculated as AFu(mi,di,ai,qi) = (qi−Qu)su aiif mi≤Mu^di≤Du 0 otherwise, (4.1) where we remind that su is the computational power of Eu . Note that Qu is equal to the current time if the queue is empty. 4.2.1. Measuring the Execution Time This policy is the first to use a key element in scheduling: the execution time of a task in a certain node. We assume that this time can be computed as ai/su . However, measuring these two values raises several problems. First, what unit should we use with them? We said in Section 3.1.1 that we measure ai in millions of FLOPs, so su is measured in millions of FLOPs per second. While this is valid for a theoretical analysis, in practice it is less useful. The task length and the computing power are not comparable between different architectures because they include instructions of varying complexity. Besides, the execution of a task is also affected by several architecture-dependent delays that are difficult to predict: cache misses, failed branch predictions, data and control dependencies, etc.
MMP: Makespan Minimization Policy 29 The second question is how to estimate ai and su . While the later may be easier to measure with benchmarks, the former has several implications. A few tasks may have a constant execution length, but most depend on the conditional branches and loops that the flow of execution follows. Furthermore, there should be an automatic method of estimating ai without actually executing a task. A possible solution would be to adapt the work by Dinda [ 40 ], who suggests estimating the execution time of a task in Eu with the nominal execution time in an unloaded, reference node, tnom , and the load of Eu . To provide a good estimation, the load of Eu is constantly monitored and predicted. Then, we could treat ai as the normalization of tnom with the inverse speed of the reference node, and predict su on every node as Dinda predicts the load. Many other works on performance prediction propose methods to estimate the node availability [ 87 , 56 , 60 ]. Task length estimation methods are covered, for instance, by works on worst-case execution-time (WCET) estimation [ 94 ]. Unfortunately, these works show that the best results are obtained with sample executions of all or part of the code to measure. For the sake of simplicity, we have used the millions of FLOPs as the task length unit to carry out the analysis and experiments of this thesis. It is realistic enough to evaluate the STaRS scheduling model and policies. Furthermore, we assume that node Eu has a constant dedicated computing power of su for the execution of remote tasks, unless it fails. In a real implementation, the previous considerations would have to be taken into account. 4.3. Global Policy 4.3.1. Availability Information So, the properties needed to compute the availability function of a node consists of its memory, disk space, computing power and end queue time. Thus, in this policy, the sampled function for a set of nodes is ( M,D,s,Q,v ), with a sample for each of the previous node parameters plus the number of nodes v . The clustering of such sampled function uses the same scheme for scalar parameters as the IBP policy explained in Section 3.2.2. We have also decided that the sum operation of f1and f2should return SUM(f1,f2) = min(f1.M,f2.M), min(f1.D,f2.D), min(f1.s,f2.s), max(f1.Q,f2.Q),f1.v+f2.v. (4.2) It calculates the minimum available memory and disk space that can be found in each node of f1 and f2 , as usual. But it also gives the minimum speed a task is going to be executed at, and the maximum time a task is going to start at. This may be seen as an excessively conservative option. Unlike the memory and disk constraints, failing to estimate the speed or the queue length is not going to avoid the allocation of a task. In this case, it would make sense to adopt a more optimistic approach with the computing power and queue end time. This can be done calculating their weighted average as avgs=f1.s f1.v+f2.s f2.v f1.v+f2.v, avgQ=f1.Q f1.v+f2.Q f2.v f1.v+f2.v,
30 MMP: Makespan Minimization Policy which would yield SUM(f1,f2) = min(f1.M,f2.M), min(f1.D,f2.D), avgs, avgQ,f1.v+f2.v. With this approach, a sampled function advertises shorter queues and faster nodes than with the previous one. So, the forwarding algorithm is able to allocate more tasks to the set of nodes of a sampled function. This results in requests needing to climb less tree levels, thus reducing the network traffic. On the other hand, sending too many tasks may also increase the makespan. After simulating both approaches under the same conditions, we have concluded that this optimistic approach always produces a larger makespan without a noticeable decrease of network traffic. So, in the experiments of Section 8we only use the conservative one. Time is relevant in this policy, since the availability function depends on Qu . So, the sampled function reported by each node must be updated when its queue end time changes. If the task queue is not empty, and tasks finish at their expected end time, Qu does not change until a new task is pushed at the end of the queue. If the task queue is empty, Qu is equal to the current time, so it changes continuously. However, if it is not updated, a routing node that finds a sampled function with Qu earlier than current time will automatically deduce that its represented queues have become empty, and update it by itself. This reduces the update frequency. 4.3.2. Forwarding Algorithm To minimize the global makespan, tasks must be sent to the nodes whose queue will remain shorter after allocating them. However, it is not enough to consider only the current branch. Routing nodes also need information about the queue end times in the rest of the tree. Without it, the routing can be performed by making all requests climb up the tree to the root node, but this solution would quickly flood the root with requests. On the other hand, routing nodes could send availability summaries to the children nodes, too. The summary sent to each child would come from the aggregation of the information obtained from the other child and the father. Then, every node would have information about all the tree. But also, each change in the availability would reach a much larger number of nodes. This solution would flood the network with availability summaries. As a compromise between these two extremes, routing nodes use the maximum queue end time in the rest of the tree in the forwarding algorithm, along with the availability information of their children. They receive it from their fathers. They send to each child the maximum between the values coming from their father’s and the other child’s subtrees. Being a maximum, its propagation is usually bounded to just a few branches, so the traffic generated is negligible. Algorithm 4.1 shows how the maximum queue end time is used in the forwarding algorithm. It starts by finding the expected makespan that will be obtained if all the tasks in the request are allocated. The availability function, with the information of the whole branch ( Ru.info ), returns the number of tasks that can be allocated for a certain makespan ( medMakespan ). The minimum expected makespan ( minMakespan ) is found with a binary search, until the availability function returns the same number of tasks the request contains. This new makespan cannot be longer than β times the longest makespan in the tree ( Ru.maxMakespan ), which is the difference between the maximum queue end time and the current time. If it is longer, and the request did not came from the father node, the request is sent to the father to look for a shorter makespan or until the root is reached. Otherwise,
MMP: Makespan Minimization Policy 31 Algorithm 4.1 Forwarding algorithm for the MMP policy. Pre: Ruis this routing node. request is the request. Post: reqLeft , reqRight and reqFather are the resulting requests to be sent to the left child, right child and father nodes, respectively. 1: procedure FORWARD(request) 2: minMakespan ←0 3: maxMakespan ←BIG_INT .Start with a very big value for the maximum. 4: while minMakespan ≤maxMakespan +1do 5: medMakespan ←(minMakespan +maxMakespan)/2 6: numTasks ←AF(Ru.info, request.PR, medMakespan) 7: if numTasks ≤request.nthen 8: minMakespan ←medMakespan 9: else 10: maxMakespan ←medMakespan 11: end if 12: end while 13: isTooMuch ←minMakespan > βRu.maxMakespan 14: if ¬ISROOT(Ru)V¬FROMFATHER(request)VisTooMuch then 15: reqFather ←request 16: else 17: numLeft ←AF(Ru.leftInfo, request.PR, minMakespan) 18: EXTRACT(request, numLeft, reqLeft) 19: reqRight ←request 20: end if 21: end procedure the availability function is used with the summary of the left child to calculate how many tasks can be sent to its branch, and the rest are sent to the right child. The β factor is a tradeoff between load balance and makespan minimization. The lower it is, the shorter the makespan will be, but the requests will usually climb more levels of the tree until they find a good set of nodes, so the upper levels will get more loaded. In Chapter 8we empirically show that 0.5 is the best value. The routing pattern of this policy is different from the one of the IBP policy. All the tasks are sent upwards until the expected makespan in the current branch gets shorter than a certain threshold. Then, they are distributed among all the nodes of that branch. In this way, the forwarding algorithm is capable of estimating a similar queue end time for all the nodes that will receive one of the tasks. Otherwise, these nodes could end up with very different queue end times, which is contrary to the objective of minimizing the makespan.
DP: A Policy for Applications with Deadlines 39 and h.ltl(δ) = f.v X i=1h.l(δ)−fi.l(δ)+ g.v X j=1h.l(δ)−gj.l(δ)= =f.vh.l(δ)−f.l(δ)+ f.v X i=1f.l(δ)−fi.l(δ)+ +g.vh.l(δ)−g.l(δ)+ g.v X j=1g.l(δ)−gj.l(δ)= =f.vh.l(δ)−f.l(δ)+g.vh.l(δ)−g.l(δ)+f.ltl(δ) + g.ltl(δ). (5.8) Once we calculate the sum of two sampled functions and obtain the MSE for all three parameters, the distance operation is again a weighted sum of their normalized values. The normalization function for the l ( δ )function is similar to the scalar case. The availability summary includes the minimum and maximum l ( δ )functions, let them be called lmin ( δ )and lmax ( δ ). Then, the normalization function is NORMl(m) = m Zh νlmax(δ)−lmin(δ)2dδ . (5.9) 5.3.2. Forwarding Algorithm The forwarding algorithm and routing pattern of the DP policy are similar to the ones in the IBP policy, shown in Algorithm 3.2 and Figure 3.5. The forwarding algorithm obtains the list of sampled functions whose nodes may fulfill the deadline of the new application. The difference is that the DP policy sorts this list having time into account, too. It selects first those sampled functions whose nodes have an available amount of computation before the application deadline closer to its task length. In this way, it tries to minimize the clearance between tasks in each execution node, so that its time is well used. As we said before, minimizing this remainder increases the probabilities of allocating future applications with longer tasks. Finally, if there are any unassigned tasks, they are sent upwards to look for execution nodes in further branches. If the root node is reached, the unassigned tasks are simply discarded because they cannot be allocated at the moment. This is interpreted by the submission node as a failure in the discovery process. Since deadlines are a quality of service, users must accept that when a task is not allocated, they must look for a longer deadline. The root node does not notify discarding a request, for two reasons. The first one is that it would generate additional traffic. The second one is that with immediate notifications, users would first send a request with a very short deadline, expecting it to fail and slowly increasing the deadline until all tasks are allocated. In this way, they would always obtain the best result for them, with an important increase of load and traffic for every one else. Then, the policy would loose its purpose. The absence of notification is solved by setting a timeout at the submission node of thirty seconds. Having to wait for this time between requests makes it very difficult for users to cheat like that.
Chapter 6. WDP: DP for Workflow Applications “All things are difficult before they are easy.” — Thomas Fuller 6.1. Scheduling of Workflow Applications In this chapter, we explore the possibility of adapting STaRS to a different application type. One of the most common types, specially in scientific domains, are the workflow applications. However, traditional grid platforms, like those based in Condor [ 90 ], show poor scaling with big workflows (up to a hundred thousand of tasks) due to scheduler overheads [ 24 ]. Scaling is even worse if most tasks in the workflows have a short duration (in the range of minutes). In [ 57 ], advance reservations, multi-level scheduling and infrastructure as a service are explored to reduce these overheads. In [ 41 ], the authors propose a distributed double-layer scheduling model for grid workflows. It contains a global scheduler that aggregates information about resource status, to avoid full knowledge. It also considers resource availability fluctuations to allocate the workflows. However, there is no study about the scalability of the system and their experiments, using 3 resource clusters with 10 nodes each, do not allow to evaluate this feature either. In [ 79 ], a scheduling algorithm is described based on the cooperation of distributed workflow brokers. A distributed hash table provides a decentralized coordination space that is responsible for the resource discovery and scheduling, and it uses a FCFS allocation strategy. So, it seems reasonable to apply a decentralized model to the scheduling of workflow applications. We present the Workflow-with-Deadlines Policy (WDP), an extension of the DP policy for workflow applications with time constraints. It includes a workflow decomposition process that allows the global scheduler to match sequences of dependent tasks with time constraints. Since this is a first approach, we only consider available computing time as the availability information of execution nodes. Likewise, the information model is simpler than the ones shown in the previous policies. 6.1.1. Workflow Management A workflow is modeled after a DAG G ( T,D ), where nodes represent tasks and edges represent dependencies between them. For each τi,τj∈T , an edge ( τi,τj ) ∈D exists if task τj depends on τi . In that case, task τimust finish before τjmay start. An example can be seen in Figure 6.1. 41
42 WDP: DP for Workflow Applications τ10 τ6 τ7 τ8 τ9 τ5 τ2 τ3 τ4 τ1 Figure 6.1.: A typical DAG with 10 tasks and several dependencies between them. We divide workflows into several parts that can be submitted concurrently, to exploit the inherent parallelism of the graph structure. Classic workflow scheduling algorithms[ 16 ]usually divide them by DAG levels. Two tasks τi and τj are in the same level when an edge ( τj,τi )or ( τi,τj )does not exist. This decomposition is suitable when the objective of the scheduler is to optimize makespan. A set of resources is reserved in advance, and levels are allocated one after another minimizing the queue length of these resources. However, we use a different approach, since we allocate resources as they are discovered. We divide the DAG into sequences of dependent tasks. A sequence is an ordered set of tasks S = {τi| 1 ≤i≤n} where ∀i, 1 ≤i<n, ( τi,τi+1 ) ∈D . The longest sequence gets its deadline directly from the DAG. Shorter ones get their time constraints from the sequences they depend on, as they are allocated. The decomposition is depicted in Algorithm 6.1. It starts by extracting the sequence that contains the longest path of the DAG, in terms of task length. Then, at each iteration, it extracts the longest sequence Sso that edges may only go: 1. From an already extracted task to the first task of S. 2. From the last task of Sto an already extracted task. Note that any DAG may be decomposed this way, as at each iteration at least a sequence of one task is eligible. Thus, all tasks are eventually assigned to a sequence. Once the DAG is decomposed into sequences, they are assigned not only a deadline, but also a startline. Every sequence must start after its startline and end before its deadline. The startline of a sequence is calculated after the task it depends on is allocated, because we know when it is supposed to finish. Likewise, its deadline is calculated after the tasks that depend on the sequence are allocated, because we know when they must start to meet their own deadline. The longest sequence is allocated first, because its deadline is provided by the user and the startline is the current time. After it is allocated, all the sequences that get their constraints from it get prepared. At each step, all the prepared sequences may be sent concurrently, as they do not depend on each other. Thus, the submission process consists of a set of stages, in which all the prepared sequences are sent. We define the width of a workflow as the number of stages needed to submit the complete workflow. Likewise, we define the minimum length of a workflow as the sum of the task lengths in the critical path, and the total length of a workflow as the sum of all the task lengths. As we show in the experimental results of Section 8.6, these properties have relevant impact in the system performance.
WDP: DP for Workflow Applications 43 Algorithm 6.1 Extract the sequences of dependent tasks from G Pre: Gis a DAG with tasks in G.Tand dependencies in G.D. Post: Scontains the sequences extracted from G. 1: function DECOMPOSITIONINTOSEQUENCES(G) 2: Extract longest path from Gto sequence S. 3: S←SSS 4: while G.T6=;do 5: Extract the tasks from G.Tthat form the longest path Sso that: • ∀si∈S|si6=s1,@τi∈S|(τi,si)∈G.D • ∀si∈S|si6=sn,@τi∈S|(si,τi)∈G.D 6: S←SSS 7: end while 8: return S 9: end function δ1δ2δ3δ4 hs1he1hs2he2 τ1τ2τ3τ4 Figure 6.2.: Task queue with four tasks, and the holes available between them. 6.2. Local Policy This policy, like the DP policy, uses an EDF-based local scheduler. For a sequence of tasks, the deadline of the sequence is applied to the last task. Then, the deadline and estimated execution time of each task is used to calculate the deadline of its predecessor. Every time the task queue changes, the availability of an execution node is recalculated, to report when and of what length new sequences could be accepted. To obtain this information, all tasks in the queue are pushed to their deadline, except the currently running task, because task preemption is not allowed. Then, a list is built up from the “holes” of availability that exist between tasks. These holes represent a simplified view of the maximum length a new sequence may have to avoid making another task miss its deadline, if it is accepted. Figure 6.2 shows an example of a task queue with four tasks. A hole that starts at hsi and ends at hei , represents the maximum length a sequence may have, with a deadline in [ δi,δi+1 ), so that in EDF order it would execute before τi+1 . Actually, if τ2 is executed just after τ1 , the hole between τ2 and τ3 would be longer, but it is not considered for sake of simplicity. Note that there is no hole between τ3and τ4, as τ4should have already started by δ3. The computational power of the node is also reported, so that it can be used to estimate how much work is performed by the node in an arbitrary period of time; for instance, when a node is completely idle.
44 WDP: DP for Workflow Applications Algorithm 6.2 Generate a set of ninterval endpoints from the current time ν. Pre: nis the number of endpoints to create, dmin is the minimum interval duration. Post: Lcontains ninterval lower endpoints. 1: procedure CREATEENDPOINTS(n) 2: L←; 3: t0←ν−BEGINNINGOFDAY(ν) 4: d1←dmin 5: for i←1to ndo 6: L←LS{dt0/die+1)×di} 7: di+1←2di 8: end for 9: end procedure 6.3. Global Policy 6.3.1. Availability Information Once the set of holes is created, it is sent to the father routing node so that it will be aggregated with the sets of the other execution nodes in that branch. However, aggregating sets of arbitrary holes would be impossible. Every pair of holes start and end at different times. For this reason, we have developed a method to approximate them to a set of fixed values. Since this is a first approach, we have simplified the availability information model. An availability summary only contains one sampled function, with the number v of nodes it represents and an ordered set of time intervals. Intervals are consecutive, so interval Ii ends at the time interval Ii+1 starts. Each one contains the list of holes that finish within its endpoints. Within each list, holes are further classified with two criteria. First, holes are grouped by their span. A hole that starts in interval Ii and finishes in interval Ii+k has a span of k +1 intervals. So, a hole that starts and finishes in the same interval has a span of one interval. However, holes with the same span may not have the same availability, as this depends on the computational power of each execution node. So, they are also classified by availability levels. Instead of using the clustering algorithm presented in Section 3.2.2, we select a suitable set of availability levels and interval endpoints that allow an easy aggregation of summaries. They must have similar values in all the availability summaries that are to be aggregated together. First, the availability of every hole is approximated to the immediately lower power of two. Like with the DP policy, this conservative method enables finding nodes where tasks meet deadlines. Also, by using powers of two, the error is always lower than 50%. The interval endpoints are generated with Algorithm 6.2. It only generates the lower endpoints; the upper endpoint of an interval is equal to the lower endpoint of its successor. Its objective is that summaries created at moments near in time have most endpoints in common. To accomplish this, it generates a set of lower endpoints whose difference with the beginning of the day is a multiple of certain durations. In the algorithm, the creation time t0 is the current time ν , relative to the beginning of the day. Then, dmin is the minimum interval duration. The first lower endpoint is t0 + dmin , rounded up to the next multiple of dmin . Then, the duration is doubled for every successive
WDP: DP for Workflow Applications 45 Table 6.1.: Interval endpoints generated from three example creation times t0 . The difference between each endpoint and t0is always between one and two times the intended interval duration. Intended interval duration t05’ 10’ 15’ 30’ 1h 2h 4h 8h 16h 1 day 4:33 4:40 4:50 5:00 5:30 6:00 8:00 12:00 16:00 1d:00:00 2d:00:00 5:48 5:55 6:00 6:15 6:30 7:00 8:00 12:00 16:00 1d:00:00 2d:00:00 17:17 17:25 17:30 17:45 18:00 19:00 20:00 1d:00:00 1d:08:00 1d:16:00 2d:00:00 interval, again to maintain the error of the approximation under 50%. So, for interval i , its duration di is dmin × 2 i−1 and it starts at t0 + di rounded up to the next multiple of di . The result is that the difference between the lower endpoint of interval i and t0 is between one and two times di . Table 6.1 shows three examples of the output of this algorithm. There, di is slightly modified to provide more “human-readable” results. Instead of being doubled at every step, it jumps from 10 minutes to 15 minutes, and from 16 hours to 1 day. 6.3.2. Forwarding Algorithm The forwarding algorithm decides how and where to route task sequences. It tries to find holes with enough availability to execute all the tasks in the sequence. Like in the DP policy, it looks for holes that better fit the sequence requirements. That is, one that starts short before the startline and ends short after the deadline of the sequence, and with an availability similar to the sequence length. In this way, it leaves bigger holes in case a longer sequence is received later, trying to improve resource usage. The forwarding algorithm of this policy is shown in Algorithm 6.3. First, GETPARTITIONS calculates all the possible partitions of the sequence into multiple consecutive subsequences. The algorithm iterates them, in ascending order of number of subsequences. At each iteration, GETHOLES looks for a set of holes that can accept each of the subsequences of the partition. They cannot overlap, to respect the dependencies between consecutive subsequences. When a hole is found for each subsequence, a new request is constructed for each one and sent to the corresponding subbranch. For each subsequence, the startline and deadline are taken from its respective hole endpoints. Finally, if the forwarding algorithm finds no hole for any of the subsequences, the original request is routed to the next level of the tree to try further.
46 WDP: DP for Workflow Applications Algorithm 6.3 Forwarding algorithm for the WDP policy. Pre: Ruis this routing node. request is the request. Post: reqLeft , reqRight and reqFather are the resulting requests to be sent to the left child, right child and father nodes, respectively. 1: procedure FORWARD(request) 2: P←GETPARTITIONS(request.sequence) 3: for numParts ←1to request.sequence.length do 4: for all S∈ { s∈P|s.length =numParts }do 5: holes ←GETHOLES(Ru.info, S) 6: if |holes| ≥ numParts then .There is a hole for every sequence. 7: reqLeft ←GETLEFTPARTS(holes) 8: reqRight ←GETRIGHTPARTS(holes) 9: return 10: end if 11: end for 12: end for 13: reqFather ←request .If we get here, there was no suitable partition. 14: end procedure
Chapter 7. FSP: Fair Share Policy “These men ask for just the same thing, fairness, and fairness only. This, so far as in my power, they, and all others, shall have.” — Abraham Lincoln 7.1. Fairness as the Scheduling Objective The Fair Share Policy (FSP) tries to provide a similar share of the platform among the scheduled applications. The fair sharing of resources has been deeply studied before in other areas, like networking [ 49 , 68 ], which consider the amount of data to be transferred by each user. When scheduling multiple applications, we consider the amount of computation that each user wants to get done. In this case, the most suited metric seems to be the maximum stretch, or slowdown [ 69 , 62 ]. The stretch of an application is defined as the ratio of its response time under the concurrent scheduling of applications to its response time when it is the only application executed on the platform. Let ri be the release time of an application, ei its finish time and tR its response time in a dedicated platform, its stretch Siis calculated as Si=ei−ri tR . (7.1) Ideally, a fair share of the platform is obtained scheduling applications so that all of them obtain the same stretch. Like it happened with the MMP policy, this is only possible in practice with offline scheduling and divisible load. We consider a more constrained environment, with online scheduling and atomic tasks, so the best tradeoff is obtained by minimizing the maximum stretch among all applications. Previously, Benoit et al. [ 14 ]have studied the minimization of maximum stretch for concurrent applications in a centralized setting. In particular, their study shows that interleaving tasks of several concurrent bag-of-tasks applications performs better than scheduling each application after the other. The main problem we face is that computing the response time of an application in a dedicated platform requires full knowledge of its characteristics. While easy to perform in a centralized context, it is unthinkable in a decentralized one. However, with some reasonable assumptions we can find a good approximation. 1. Applications have much less tasks than nodes in the platform. 47
48 FSP: Fair Share Policy 2. The distribution of computing power changes very little. We already presented a first design of this policy in [ 37 ]. In this chapter, we refine its methods to obtain better results. 7.1.1. Slowness The first assumption is that applications have much less tasks than nodes in the platform. From the traces used in the experiments of Section 8, the number of tasks per application in a common cluster computing environment is distributed as: •The application with more tasks has 1120 tasks. •95% of applications have less than 128 tasks. •75% of applications have less than 16 tasks. •50% of applications have only one task. Since we aim for platforms with a hundred thousand nodes or more, to minimize tR of an application with ni tasks, we have to allocate a task to each of the ni fastest nodes. Then, the response time would be the time needed by the slowest of these nodes to execute its task. The second assumption is that the distribution of computing power changes very little. Even more, we assume that the fastest nodes will have a similar computing power. This is more complicated, since the computing power distribution is usually skewed towards the lower values, but with the relation of a thousand tasks to a hundred thousand nodes, we think it is reasonable. In that case, we can approximate tR as ai/smax , where smax is the computing power of the fastest node. In practice, this value is still unknown. But since we assume that it changes very seldom, we turn the problem of minimizing the maximum stretch into minimizing the maximum ratio between the stretch and smax . This ratio ziis zi=Si smax =ei−ri ai . (7.2) Note that zi is the inverse of the effective speed at which a task has been executed, so we call it slowness. 7.2. Local Policy 7.2.1. Minimizing the Local Maximum Slowness The local scheduler only knows about the applications which have at least one task allocated to its execution node. So, its objective is to minimize the maximum slowness among these applications. Moreover, in order to estimate the eventual slowness of an application, it must calculate its end time ei based only on the tasks that are contained in the queue. Therefore, two local schedulers may compute a different slowness for the same application. The global scheduling policy is in charge of minimizing this unbalance.
FSP: Fair Share Policy 55 In practice, we have observed that ignoring t yields better results. Since we are minimizing the maximum slowness, we do not want to get a higher real value than the one we estimated. Besides, only with the piece endpoints and parameters, we have a very limited information about Eu queue. First, in those pieces where we subtract t from uj , we do not know whether τk is still setting the maximum slowness. Second, the endpoints of those pieces may change in an unknown way, yielding a wrong estimation in the adjacent ones. Nevertheless, note that the previous considerations assume that the queue does not change. So, in any case, zu ( a )pieces must be recalculated every time a new task is added to the queue or an old one finishes. Estimating zu(a)for more than one task This is how Equations 7.6 and 7.7 look like if we allocate ntasks of size ak: zk= nak+Pk−1 j=1aj+ (rk−rk)su aksu =Pk−1 j=1aj aksu +n su zi>k=nak aisu +Pi j=1,j6=kaj+ (rk−ri)su aisu Once when we have the piece endpoints and parameters of zu ( a )for one task, we can safely say that, multiplying the vj and wj parameters of each piece by n , we obtain a good estimation of zu ( a )for n tasks. Two problems arise, similarly to the estimation with different ri . First, in pieces where τk was setting the maximum slowness, a task later in the queue may take over, or vice versa. Second, the piece endpoints may also change, like in the previous case. Nevertheless, the experiments have shown that, in this case, the estimation is better than not taking action. 7.3. Global Policy 7.3.1. Availability Information Management Similarly to the DP policy, the sampled function of the FSP policy consists of the parameters ( M,D,Z,v ). It contains the number of nodes v , a sample M and D for the available memory and disk space, and a set of samples Z = { ( αj,uj,vj,wj ) , 1 ≤j} . These samples define the pieces of the z ( a )function of the nodes represented by this sampled function, with their endpoints and parameters. Only the left endpoint αj is needed because pieces are adjacent, the right endpoint is the next piece’s left endpoint. The last piece has no right endpoint. The first two parameters of a sampled function are scalar values, as usual. Z is a functional parameter, and all the considerations that applied to L in Section 5.3.1 apply to Z now. In particular, for notation, the function represented by f.Z is written as f.z ( a ). As we said in the previous section, zu ( a )must be recomputed in Eu every time its queue changes, so an availability summary is sent to the father routing node with this frequency.
56 FSP: Fair Share Policy Figure 7.4.: Two examples of how to join two pieces (black) into one (red). Sum operation The sum operation is computed as SUM(f1,f2) = min(f1.M,f2.M), min(f1.D,f2.D), max(f1.Z,f2.Z),f1.v+f2.v. (7.9) As we said before, since we are trying to minimize the maximum slowness, we prefer that the estimation overshoots the real value. So, the sum operation calculates the maximum between f1.Z and f2.Z . Like with the L parameter of the DP policy, the result is a set of pieces whose endpoints come from the endpoints of f1.z ( a )and f2.z ( a ), and from the points where they cross each other. So, the number of pieces in SUM ( f1,f2 ) .Z may be up to three times the number of pieces of f1.Z or f2.Z . In this case, we reduce its number by joining two consecutive pieces into a new one. It covers the same range as the pieces it replaces, and its parameters are calculated so that it maintains the same value at its endpoints, and provides a higher or equal estimation of the slowness between them. Figure 7.4 show two examples of this process. Joining two pieces carries an error in the estimation, so the sum operation iteratively joins those pieces that minimize the error, until a fixed number of pieces is reached again. Then, these estimation errors are also taken into account in the distance operation. Distance operation In the distance operation of the FSP policy, the parameters M and D are treated as scalar values, as explained in Section 3.2.2, while parameter Z is treated like the parameter L of the DP policy in Section 5.3.1. Let f.z ( a )be the function described by the pieces in f.Z . Then, the sum operation of fand gwould obtain hwith the MSE of parameter Z, h.msez=1 h.v f.v Zh νh.z(a)−f.z(a)2da +f.msez!+ +g.v Zh νh.z(a)−g.z(a)2da +g.msez!+ +2Zh νh.z(a)−f.z(a)f.ltz(a)da + +2Zh νh.z(a)−g.z(a)g.ltz(a)da , (7.10) and the linear term of parameter Z, h.ltz(a) = f.vh.z(a)−f.z(a)+g.vh.z(a)−g.z(a)+f.ltz(a) + g.ltz(a). (7.11)
FSP: Fair Share Policy 57 For the normalization function of the Z parameter, the availability summary includes the minimum and maximum z ( a )functions as usual. Let them be called zmin ( a )and zmax ( a ), the normalization function is NORMz(m) = m Zh νzmax(a)−zmin(a)2da . (7.12) 7.3.2. Forwarding Algorithm The FSP policy tries to minimize a global property, the maximum slowness. So, its forwarding algorithm and routing pattern are closer to those of the MMP policy. Tasks must be sent to those nodes whose maximum slowness will remain lower after accepting them. To have an idea about the global maximum slowness, besides the availability information that comes from each child, routing nodes receive the maximum slowness of the rest of the tree from their father. Like in the MMP policy, being a maximum, it seldom changes and its impact in the traffic is negligible. Requests climb up the tree until a routing node decides that, with the nodes of its branch, the maximum slowness will be no greater than β times the maximum slowness of the rest of the tree. Then, it divides the set of tasks of the request between its children. In Chapter 8we also show that, in this case, the value of β that obtains the best results is 0.04. The forwarding algorithm of the FSP policy appears in Algorithm 7.6. First, it calculates the minimum slowness that can be reached allocating the tasks of the request to the nodes in the current branch. This is done by function GETMINSLOWNESS, in Algorithm 7.7. It creates a list candidates with all the sampled functions and estimates the slowness obtained by assigning one task to the nodes of each candidate. Procedure PURGE eliminates from the list the worst candidates so that there remains just enough nodes to allocate request.n tasks. At this point, the partial result is the maximum slowness among them, in candidates.last.slowness . However, if we assign more tasks to the nodes with lowest slowness, they could still get a lower slowness than the other nodes. To find this out, the algorithm iteratively checks whether it obtains a better result by assigning one task more to some sampled function. At iteration i , the sampled functions that obtained i− 1 tasks in the previous one are tested with one task more. This is done with function ESTIMATEZ, that estimates the value of function z ( a )with i tasks. If they still get a lower slowness than the partial result, they get that new one task, and the list of candidates is purged again. The process stops when it cannot obtain a lower slowness. The forwarding algorithm then compares the obtained slowness with β times the maximum slowness in the rest of the tree. If it is lower, the availability function decides how many tasks to send to each child. It is easy to see that the routing pattern generated by this forwarding algorithm is the same as that of the MMP policy.
58 FSP: Fair Share Policy Algorithm 7.6 Forwarding algorithm for the FSP policy. Pre: Ruis this routing node. request is the request. Post: reqLeft , reqRight and reqFather are the resulting requests to be sent to the left child, right child and father nodes, respectively. 1: procedure FORWARD(request) 2: minSlowness ←GETMINSLOWNESS(Ru.info, request) 3: isTooMuch ←minSlowness > βRu.maxSlowness 4: if ¬ISROOT(Ru)V¬FROMFATHER(request)VisTooMuch then 5: reqFather ←request 6: else 7: numLeft ←AF(Ru.leftInfo, request.PR, minSlowness) 8: EXTRACT(request, numLeft, reqLeft) 9: reqRight ←request 10: end if 11: end procedure Algorithm 7.7 Get the minimum slowness that can be reached when assigning the tasks in request to the nodes described by info. Pre: info is an availability summary. request is the request. Post: candidates.last.slowness is the minimum slowness that can be reached when assigning the tasks in request to the nodes described by info. 1: function GETMINSLOWNESS(info, request) 2: candidates ←; 3: for all sf ∈info do 4: candidates ←candidates S{(sf,sf.z(info.a), 1)} 5: end for 6: PURGE(candidates, request.n) 7: oneMoreTask ←true 8: i←1 9: while oneMoreTask do 10: i←i+1 11: oneMoreTask ←false .If nothing happens, this is the last iteration. 12: for all (sf, slowness, n)∈candidates |n=i−1do 13: if ESTIMATEZ(info.a,i)<candidates.last.slowness then 14: candidates ←candidates \ {(sf, slowness, n)} 15: candidates ←candidates S{(sf, ESTIMATEZ(info.a,i),i)} 16: oneMoreTask ←true 17: end if 18: end for 19: PURGE(candidates, request.n) 20: end while 21: return candidates.last.slowness 22: end function
Chapter 8. Experimentation “It doesn’t matter how beautiful your theory is, it doesn’t matter how smart you are. If it doesn’t agree with experiment, it’s wrong.” — Richard P. Feynman We have measured the scalability, fault-tolerance and performance of our proposal through a set of tests and simulations. We have first developed tests to evaluate the accuracy of the aggregation scheme. They are run by a specific evaluation program that aggregates the information of a set of nodes in the same way that would be done in the tree. Then, we have also developed an ad-hoc discrete event simulator (DES) to observe our model under more realistic conditions. This DES is written in C++, focused on minimizing the memory footprint to simulate as many nodes as possible. Its details can be found in Appendix B. The five presented policies have been tested, but only the IBP, MMP, DP and FSP policies are compared with each other. Due to its characteristics and it being a first attempt to schedule a different kind of application, the WDP policy results are presented in Section 8.6. 8.1. Aggregation Tests We have studied the accuracy of the aggregation scheme for the different kind of parameters that have appeared in this thesis. To evaluate its scalability, we have varied the system size and the number of sampled functions per summary. We understand the accuracy of the aggregation scheme, applied to a set of nodes, as the fraction of their actual availability that is represented in the resulting summary. With 100% accuracy, the resulting summary perfectly represents the actual availability. The 0% accuracy corresponds to the minimum value that the scheme could potentially calculate. For instance, to aggregate the available memory of two nodes, we compute their minimum. So, the minimum of all the nodes would be an accuracy of 0%. Note that this is more restrictive than assigning a 0% accuracy to no available memory. We have developed a specific evaluation program that calculates this fraction. The parameters of the program are the size of the network ( N ), the policy being tested and the maximum number of sampled functions per summary ( SFmax ). First, it generates the availability information of a set of N nodes, setting their properties and queue states at random with a uniform distribution. This distribution of values is the worst case for a clustering algorithm. Then, it aggregates this information recursively, imitating the organization of the nodes in a balanced binary tree. Finally, it compares the result with the actual availability of all the nodes and calculates the accuracy of the aggregation. 59
60 Experimentation 50 100 150 200 250 300 0 25 50 75 100 10 (a) Accuracy (%) Memory Disk space 50 100 150 200 250 300 0 25 50 75 100 10 (b) Accuracy (%) Memory Disk space Processor speed Queue length 50 100 150 200 250 300 0 25 50 75 100 10 (c) Accuracy (%) Memory Disk space Available FLOPs 50 100 150 200 250 300 0 25 50 75 100 10 (d) Maximum sampled functions per summary Accuracy (%) Memory Disk space Slowness Figure 8.1.: Aggregation accuracy of a set of 1024 nodes for an increasing SFmax , for the (a) IBP, (b) MMP, (c) DP and (d) FSP policy parameters. The accuracy represents the fraction of the actual availability of a set of nodes that is represented in the aggregated summary. It is very similar for the scalar parameters.
Experimentation 61 128 256 512 1K 2K 4K 8K 16K 32K 64K 128K 0 25 50 75 100 (a) Accuracy (%) Memory Disk space 128 256 512 1K 2K 4K 8K 16K 32K 64K 128K 0 25 50 75 100 (b) Accuracy (%) Memory Disk space Processor speed Queue length 128 256 512 1K 2K 4K 8K 16K 32K 64K 128K 0 25 50 75 100 (c) Accuracy (%) Memory Disk space Available FLOPs 128 256 512 1K 2K 4K 8K 16K 32K 64K 128K 0 25 50 75 100 (d) Number of nodes Accuracy (%) Memory Disk space Slowness Figure 8.2.: Aggregation accuracy with 200 sampled functions per summary for an increasing number of nodes, for the (a) IBP, (b) MMP, (c) DP and (d) FSP policy parameters. The accuracy represents the fraction of the actual availability of a set of nodes that is represented in the aggregated summary. It is very similar for the scalar parameters.
62 Experimentation We execute the program for each policy and several network sizes and SFmax values. At each step, we double the number of nodes to increment the tree height one more level. Figure 8.1 shows the aggregation accuracy with different SFmax , for the four policies and a set of 1024 nodes, which would be found at level 10 of the tree. It can be seen that the accuracy quickly improves when we increase SFmax in the low value range. It hardly improves anymore when we increase SFmax over 200 sampled functions per summary. We show later how SFmax affects the network traffic and each policy performance, in order to choose an appropriate value for each situation. Then, Figure 8.2 presents the aggregation accuracy with an increasing number of nodes, for an SFmax of 200 and each policy. The decrease of accuracy is very low compared to the increase in the number of nodes. This figure also highlights that it is difficult to normalize the MSE of the functional parameters. The available amount of FLOPs before deadline of the DP policy gets noticeably less accuracy than the memory and disk space parameters, while the slowness per task length of the FSP policy gets more, because they are not correctly weighted in the distance operator. But in general, the clustering algorithm provides a well balanced accuracy among the different resource types, even when their values lay in very different intervals. 8.2. Simulation Setup Besides the aggregation tests, we have also developed a DES that executes our scheduling model in a network of nodes for a certain amount of time. The network size is configurable, and there is direct communication between every pair of nodes. Nodes can send messages to each other, or messages can be injected to emulate the user behavior. To check the scalability of our model, we have simulated networks from five thousand up to a hundred thousand nodes. With the IBP policy, whose complexity is lower, we are able to simulate networks of up to a million nodes. For every message, the DES takes into account both the transmission and computational times. The transmission time of a message is calculated by modeling the end-to-end link. We test our proposal with the end-to-end link of a typical volunteer computing platform, with 10 Mbps bandwidth and a delay between 50 ms and 300 ms, and of a fast cluster interconnection network, with 1 Gbps bandwidth and a delay between 0.1 ms and 1 ms. The delay follows a Pareto distribution, as suggested in [ 97 ]. The processing time of a message is the time needed by the simulation machine to process it. It could be scaled to the computing power of each node, but we decided to not do it. In this way, we are able to compare the processing time of a request in our model with the processing time in a single centralized machine. Only task execution time is calculated with the computing power of the execution nodes. Simulations consist in having a user at each node, continuously submitting new jobs. For the generation of this workload, we have implemented the site-level simulation model proposed by Shmueli and Feitelson [ 86 , 44 ]. It simulates the interaction between users and the system to calculate the timing and parameters of each submitted application. It was designed analyzing the influence of the scheduling performance on user decisions in several production system traces. These traces contain many years of activity of about 2000 users, with information of nearly half a million jobs. The DES also uses them to extract the distribution of application parameters and node properties (power, memory and disk), so we consider it realistic enough.
Experimentation 63 25K 50K 75K 100K 0 5 10 15 20 5K Network size (Knodes) Throughput (tasks/s) IBP MMP DP FSP Figure 8.3.: Average throughput of each policy by network size. It depends on the workload characteristics, like the distribution of the number of tasks and release time. A problem we faced with Shmueli and Feitelson’s model is that the probability of a user having a break increases with the job turnaround time. Poorer performance usually yields longer job turnaround time. So, for certain performance metrics, like number of successfully finished tasks, the difference between two simulations’ results may be larger than it should just because users are breaking more often and sending less jobs in the least performing one. For this reason, instead of applying the site-level model to every simulation as is, we apply it just once, generate our own trace of workload and replay it in the simulations we want to compare together. Then, users always submit the same amount of jobs and differences in the results are only due to our design. Finally, we have developed a centralized version of each policy but the WDP. It is a centralized online scheduler that uses the same heuristics as its decentralized counterpart, but with full knowledge of the system state. In this way, it provides a reference to evaluate the advantages of a decentralized design against a centralized one. Their implementation and details can be found in Appendix C. 8.3. Scalability Results We first show how our model preserves its scalable properties regardless of the policy being used through three metrics: the throughput, the allocation time and the bandwidth usage. Figure 8.3 presents the average throughput of each policy for increasing number of nodes. Parameters other than the network size are kept constant. The absolute values are not relevant, because they depend on the workload characteristics of the site-level simulation model, and no policy tries to maximize the throughput. Instead, it is interesting to see that, for every policy, the throughput linearly increases with the system size. We define the allocation time as the time elapsed between a request is submitted and all its tasks get accepted. It measures the cost of the task allocation algorithm as perceived by the user. Figures 8.4 and 8.5 compare the average allocation time of a request with a thousand tasks between the decentralized and centralized versions of all four policies, respectively. Figures 8.4a and 8.5a show the results of the decentralized and centralized versions of each policy in the slow network link model, while Figures 8.4b and 8.5b are plotted with the results of the fast one. As predicted in Section 3.3, in the decentralized case, the increase of allocation time with network size is near logarithmic, which
64 Experimentation 5K 25K 50K 75K 100K 1s 1.5s 2s 2.5s 3s 3.5s (a) Allocation time IBP MMP DP FSP 5K 25K 50K 75K 100K 3ms 4ms 5ms 6ms 7ms 8ms (b) Network size (Knodes) Allocation time IBP MMP DP FSP Figure 8.4.: Average allocation time against network size and policy, of a 1000 task request, with the (a) slow and (b) fast network link, by the decentralized versions of each policy. 25K 50K 75K 100K 80ms 100ms 120ms 140ms 5K (a) Allocation time IBP MMP DP FSP 25K 50K 75K 100K 5ms 20ms 35ms 50ms 5K (b) Network size (Knodes) Allocation time IBP MMP DP FSP Figure 8.5.: Average allocation time against network size and policy, of a 1000 task request, with the (a) slow and (b) fast network link, by the centralized versions of each policy.
Experimentation 71 0 10 20 30 40 50 60 0 3 6 9 12 15 Simulation time (hours) Maximum slowness Centralized SFmax =200 SFmax =50 SFmax =20 Figure 8.11.: Maximum slowness among coexisting applications during the simulation, by the FSP policy, on a hundred thousand nodes, for different values of SFmax, compared to its centralized version. 32.1 and with 100 Bps it was 58.4. We conclude that maintaining the availability information up to date is determinant for this policy. With these differences, we have decided to do the rest of the simulations with the fast link model and 10 KBps of update bandwidth limit. Like in the MMP policy, we also have to set the β parameter to find a tradeoff between performance and traffic. In this case, the β parameter decides if the slowness obtained by the forwarding algorithm in a branch is short enough (see Section 7.3.2). We have tested the performance of the FSP policy for several values of β between 0 . 01 and 2. The maximum slowness is regular with β < 1, quickly increasing after this value. However, the number of requests that reach the root start growing significantly after β =0 . 04. So, we think that β =0 . 04 is the correct value for the following performance tests. Finally, we measure the performance of the FSP policy by the maximum slowness among the applications that coexist in the system. Figure 8.11 shows how it evolves during the simulation on a hundred thousand nodes. We have plotted the results for a SFmax value of 20, 50 and 200 sampled functions per summary and for the centralized version. The random allocation, on the other hand, is off the charts. Its maximum slowness quickly grows to 138, remaining there the rest of the simulation. For the decentralized version, it can be seen that it presents a noticeable variation of performance among different values of SFmax . The performance with SFmax =200 is up to 20 times that with SFmax =20. However, we are still far away from the performance of the centralized version. After 5 hours of simulation, when it was nearest to the performance of the decentralized version with SFmax =200, it was still about 40 times better. 8.5. Fault-tolerance Results To test the fault-tolerance of our model, we also perform simulations in which nodes fail. As explained in Section 3.1.4, the consequences of a node failing are that all the tasks in its queue are aborted and the availability information of its branch is lost. The information about issued application requests and tasks waiting to be finished is saved to a database. We model failures as in [ 83 ]: every time a node leaves another one replaces it, so that the system maintains its size. We assume that the tree overlay
72 Experimentation 0.47 60’ 0.36 30’ 0.264 15’ 0.168 5’ 0 0.1 0.2 0.3 0.4 0.5 0.6 (a) 0.388 60’ 0.324 30’ 0.274 15’ 0.155 5’ 0 0.1 0.2 0.3 0.4 0.5 0.6 (b) 0.243 60’ 0.193 30’ 0.160 15’ 0.076 5’ 0 0.1 0.2 0.3 0.4 0.5 0.6 (c) 0.387 60’ 0.284 30’ 0.151 15’ 0.038 5’ 0 0.1 0.2 0.3 0.4 0.5 0.6 (d) Figure 8.12.: Finished computation for simulations with churn, with median session times of 60, 30, 15 and 5 minutes, for the (a) IBP, (b) MMP, (c) DP and (d) FSP policies. It is normalized to the results without churn. is one of the works cited in Section 3.1.4 and that it recovers by itself, but we do not consider its recovery overhead. That would depend on the actual implementation used, so we focus on the cost of recovering the components of our model. Then, the system automatically recovers by redistributing the availability information and resubmitting aborted tasks. We perform these failure tests with the same parameters as the performance tests, and an SFmax value of 200. We simulate churn and catastrophic failures. Churn is the continuous process of node arrival and departure, very common in desktop grids and P2P computing platforms. In managed environments, like clusters and data centers, churn is very light. Rhea et al. [ 83 ]characterize churn by the median session time of a node. From the observed session times in various P2P systems, they suggest median session times from 5 to 60 minutes. We show how the computation finished by each policy degrades with churn in Figure 8.12, for median session times of 60, 30, 15 and 5 minutes. Values are normalized to the finished computation by each policy without churn. We understand a catastrophic failure as the simultaneous failure of a large set of nodes. It is rare, but it is most significant in managed environments; e.g. due to a power loss in a data center, a cutting in an interconnection network or a misconfiguration propagated to several virtual instances. We tested our model with all four policies when a catastrophic failure occurs after 1 day of simulation, with 5%, 10%, 20% and 40% of the nodes failing. As soon as the underlying tree structure recovers from the failure, our model is able to quickly recreate the lost availability information and resubmit the aborted tasks, so the impact in global performance of each policy is minimal. The only noticeable effect is that the recovered applications obtain a longer response time, and in the DP policy some may not meet their deadline.
Experimentation 73 τ10 τ6 τ7 τ8 τ9 τ5 τ2 τ3 τ4 τ1 Figure 8.13.: A Fork-Join graph model with 10 tasks. τ1τ2τ3 τ4τ5τ6 τ7τ8τ9 Figure 8.14.: A Laplace equation solver graph model with 9 tasks. 8.6. WDP Policy Tests The WDP policy is different to the other policies in many aspects. It is a first approach to schedule a different kind of application. So, the site-level simulation model cannot be applied to this policy, because it is not designed to generate workflow applications. Also, the availability model has been simplified, and there is nothing like a sampled function in this policy. For these reasons, and to assess whether it is worth further development of this policy, we only present results on allocation time, speed-up and computational and network costs of our decentralized proposal. Additional experiments will be carried out in the future. 8.6.1. Simulation Setup We have simulated this policy on the slow network model with one million nodes. As workflows, two different kinds of DAGs has been selected to get results on different sizes and workflow widths. The selected models are a Fork-Join and a Laplace equation solver-like graph, shown in Figures 8.13 and 8.14. The arrival of new workflows follows a Poisson process, with a mean interarrival time of 3 milliseconds. The mean arrival rate is twice the aggregated speed of all the platform, so it maintains every node busy at almost any time. Finally, the deadline of each workflow is calculated through the workflow relative priority Wj . It is the ratio between the time the submission node would need to execute the workflow and the time remaining until deadline. A value of 1 or less would mean that the sender could execute that workflow by itself, so higher values are tested.
74 Experimentation 0 2 4 6 8 10 12 0 0.2 0.4 0.6 0.8 1 Allocation time (seconds) Cumulative fraction Fork-Join 10 tasks Laplace 9 tasks Fork-Join 100 tasks Laplace 100 tasks Figure 8.15.: Allocation time for different workflow widths, in a network of 1 million nodes. 8.6.2. Results Probably, the system property that users firstly notice is response time. It depends on two factors: workflow allocation time and workflow execution time. The time a workflow needs to be allocated is directly related to its width. Figure 8.15 shows the cumulative distribution of the allocation time for both models and different sizes. As expected, the allocation time of fork-join workflows varies very little with their size, while it grows significantly for Laplace-like models, whose concurrency is lower. As with the other policies, the allocation is done in just a few seconds. The workflow execution time is the time lapse between a request is sent and the last task is finished. We define the speed-up as the ratio between the time the submission node would need to execute the workflow by itself and the actual execution time. The ratio between workflow total length and minimum length provides the maximum speed-up that can be expected in execution nodes with similar computing power as the submission node. Table 8.2 presents average speed-up values registered for different workflow relative priorities. As it can be seen, the fork-join model has an expected maximum speed-up of 3.3 while the Laplace model could get a speed-up of just 1.8. Results show that for the Laplace model, the system performs better than expected, as tasks may run in faster machines than the client’s one. However, in the fork-join model, the WDP policy is unable to reach the maximum expected speed-up. In both cases, the speed-up increases with the workflow priority as deadlines become tighter. Finally, we present the network traffic of the WDP policy. It has been tested with links of 1Mbps of bandwidth and an update bandwidth limit of 10 KBps. The average traffic along the simulation is very low, so we have studied the peaks of traffic. Like with the other policies, we measure the traffic at intervals of one second, and register the maximum of all intervals. The average peak of outgoing traffic among all the nodes was 216 Bps, 75% of the nodes only generated less than 2500 Bps of peak traffic, and the maximum peak was 8200 Bps. Likewise, the average peak of incoming traffic was 190 Bps, about 75% of the nodes received less than 1000 Bps of peak traffic, just 1% of the nodes received
Experimentation 75 Table 8.2.: Average speedup by workflow model and priority. Maximum Speed-up Workflow relative priority Wj Model 1.1 1.2 1.3 1.4 Laplace 1.8 2.05 2.09 2.18 2.22 Fork-Join 3.3 2.12 2.85 2.90 2.92 more than 8000 Bps, and the maximum value registered during the tests was 29000 Bps. Again, the impact of the communications in the user experience is very limited. 8.7. Comparison with Other Works In this chapter we have quantitatively compared the scalability and performance of our decentralized model with a centralized one, which we have also implemented. But it would also be interesting to compare ourselves with some of the work presented in Chapter 2. However, this is often a difficult task, because each one uses different metrics and parameters. Many decisions made by the authors affect the results, like the policies being used, the generation of workload or the network model. So, we have made a more qualitative comparison. After studying the characteristics that the experiments of the related work have in common with us, we have selected four of them as indicators of scalability, fault-tolerance and versatility: The number of simulated nodes, to provide an idea of how far the authors pushed their implementations; the existence of a failure scenario; the implemented policies; and the workload generation, to evaluate if the simulations were realistic enough. Table 8.3 shows the results. Decentralized resource discovery platforms offer the best numbers, but they obviously implement no scheduling policy. NodeWiz, SWORD and the work by Cardosa and Chandra show simulations of 10,000 nodes. NodeWiz authors even test a prototype implementation on 1,000 emulated nodes and 100 real nodes, and take failures into account. No other work shows results of tests with failing nodes. All of them use real data collected from platforms like PlanetLab [ 72 ]. The exception is Cohesion, whose authors only test platforms of 64 nodes with synthetic data. The works on grid schedulers with aggregated information present the tests with the least number of nodes. Their scalability is limited because they usually only consider one broker per grid domain, and they use a centralized scheduler (Kokkinos and Varvarigos) or expect every domain to know each other (Brunner et al. and Rodero et al.). Surprisingly, Rahman et al. claim to decentralize their scheduling by using a DHT to index domain capabilities, but they only test it on 100 domains. On the other hand, they test several scheduling policies, including scientific workflows. Rodero et al. and Kokkinos and Varvarigos also use traces from the Grid Workloads Archive [91]. Finally, we consider that we have improved the state of the art in relation to other decentralized scheduling platforms. They experiment with less than 10,000 nodes, no node failures, only one policy and synthetic workloads. WaveGrid does not take node failures into account, but it migrates tasks when users reclaim their nodes. Kwan and Mupala mention the resilience of their unstructured network, but it is not tested. Kim et al. compare the performance of its model with an idealistic
76 Experimentation Table 8.3.: Comparison of several distributed scheduling projects. Project N. nodes Failure Policies Workload Our model 1,000,000 Yes Various Traces NodeWiz [12]10,000 Yes N.A. Mixed SWORD [7]10,000 No N.A. Mixed Cardosa and Chandra [26]10,000 No N.A. Traces Cai and Wang [23]8,192 No N.A. Traces WaveGrid [98]5,000 No IBP Synthetic Kwan and Mupala [61]4,000 No IBP Poisson p. Kim et al. [58]1,000 No MMP Poisson p. Kokkinos and Varvarigos [59]1,000 No Various Traces Brunner et al. [18]1,000 dom. No Workflows Synthetic Rahman et al. [75]100 No Workflows Poisson p. Cohesion [85]64 No N.A. Synthetic Diet [27]50 No IBP Synthetic Rodero et al. [84]18 No various Traces centralized scheduler, like we do, with a MMP-like policy. Yet, we consider that 1,000 tested nodes are too few nowadays. Diet authors propose an extensible architecture, with plug-in policies, but only one is tested. Although only 50 nodes are tested, they perform tests on a full implementation, not on simulation. All these works generate their own synthetic workloads, usually with a Poisson process. They could present more realistic results with workloads coming from real system traces. In contrast, we test up to one million nodes, millions of tasks generated from real traces, failures of varying size and frequency, and five different policies. 8.8. Discussion Returning to the properties we considered in Section 1.2 a distributed scheduling model should have, we can evaluate now if we accomplished our goals. The model should be scalable, in the sense that it should be able to deal with an increase in the number of nodes without a noticeable impact in its performance. Our scalability results in Figures 8.4 and 8.5 show that the allocation cost has a logarithmic-like behavior with the system size, due to the concurrent allocation in different branches. We emphasize that our model can be faster than a centralized implementation in low delay interconnection networks, like those found in cloud computing facilities and data centers. Yet, it still manages to allocate a thousand tasks among a hundred thousand nodes with high link delay in less than 3.5 seconds. With its logarithmic behavior, we can extrapolate an average allocation time of less than 5 seconds with ten million nodes. Besides,
Experimentation 77 the communication overhead seems to be bounded even with an increase of the number of nodes. The bandwidth usage with the fast link model is negligible, and even the peaks with the slow link model can be considered low. This is due to the aggregation scheme, that effectively limits the traffic in the top levels of the tree. Finally, the throughput increases linearly with the number of nodes, without notice of deceleration. We consider reasonable to extrapolate these results and think that our model would behave similarly in higher scales than tested. The model should also be fault-tolerant, degrading its performance but continuing its work in the case of failure. The tests on churn and catastrophic failures confirm that our model supports high rates of failures without loosing functionality, just degrading the performance accordingly. However, churn rises some special considerations. It can be seen that churn quickly affects the performance of the DP policy. This is because deadlines are a heavy requirement and resubmitted tasks cannot meet them. On the other hand, short session times badly affect the FSP policy, while it should be expected to behave like the MMP policy. Since the FSP policy has been recently developed, we have to further investigate the reason of this problem. In general, complex policies or policies with strong requirements should get worst results under churn. Additionally, most long tasks cannot finish with the shortest session times. To overcome this problem, users should also implement checkpointing or replication into their applications, to avoid resending failed tasks too often. Finally, the model should be extensible to several different policies. By design, our generic availability information representation, tunable aggregation scheme and task routing approach make it feasible. With them, we have implemented five different policies of increasing complexity. We want to highlight that the most common IBP, MMP and DP policies obtain very good performance results compared to their centralized versions, while using availability information of very different level of detail. The WDP and FSP policies, however, still have room for improvement. The novelty of the former and the simplifications we made in the later yield moderate results. It is also worth bearing in mind that the MMP and FSP policies require a higher update bandwidth. The most probable cause is their routing pattern, but we must study this problem further.
Chapter 9. Conclusions “Yet in all those cases I finally steeled myself to seize the opportunity, and find a way to muddle through and eventually conclude that I had, in fact, chosen the right path, as risky as it seemed at the time.” — Vinton Cerf In this thesis we presented a distributed scheduling model for large-scale platforms. It aggregates availability information about execution nodes on a hierarchical overlay. Then, using that information, it forwards tasks towards the most suitable execution nodes. We claim that it reaches scales of millions of nodes, tolerates high rates of failures and supports policies with very different objectives. We provide results from trace-driven simulation tests on a network of up to a million nodes. The scalability is achieved through two main mechanisms. First, we propose an aggregation scheme that provides enough availability information to the top levels of the hierarchy without flooding them. An agglomerative clustering algorithm summarizes the availability information to avoid it growing without limit. Then, the update bandwidth is also bounded to reduce the network traffic. This scheme finds a good tradeoff between resource usage and accuracy. Second, we take a task-routing approach to scheduling. The nodes of the hierarchy use a forwarding algorithm that looks for execution nodes in several branches of the hierarchy concurrently, because they are independent. In our tests, the communication overhead is bounded and the allocation cost shows an almost logarithmic behavior with the system size. In networks of 100,000 nodes, link bandwidth of 1Gbps and link delays of under 1 millisecond, our decentralized scheduler can allocate tasks up to 10 times faster than its centralized equivalent. For instance, the decentralized version of the DP policy allocates a thousand tasks in 5ms, against 50ms of the centralized version. Meanwhile, the maximum bandwidth usage peak was under 0.25% of the link bandwidth. With delays of up to 300 milliseconds, the slowest policy still needs less than 3.5 seconds. The result is that our model is very well suited even for applications with thousands of short tasks, like many-task computing and map-reduce applications. Faults are managed with a best-effort strategy. When a node failure is detected, the tasks it was executing are resubmitted by their owners and the availability information it managed is rebuilt by its neighbors. In this way, our model is able to degrade its performance accordingly and recover its functionality. To show it, we have performed tests of churn and catastrophic failures. Even with nodes failing with a median period of only 5 minutes, our scheduler is able to continue giving a degraded service. Meanwhile, it is able to recover from the failure of an important fraction of the nodes, as long as the underlying hierarchical overlay supports it. 79
80 Conclusions Finally, we have designed our distributed scheduling model with extensibility in mind. The availability information model supports a generic set of operations to aggregate different properties of the execution nodes. Besides, the forwarding algorithm provides a common prototype to different realizations. So, we have implemented five policies, by specializing their own availability representation and forwarding algorithm. The IBP policy allocates bag-of-tasks applications to idle nodes as they become ready, if they meet the memory and disk space requirements. The MMP policy also considers queue length to minimize the global makespan. The DP policy allows the use of time constraints to schedule first the applications with a shorter deadline. The FSP policy introduces the concept of slowness to provide a fair share of the platform to every application. And the WDP policy explores the extension of the DP policy to a different application type, the workflow, that includes additional dependencies between tasks. We propose a common method to represent and clusterize scalar parameters of the execution node availability, like available memory and disk space. It is applied in the four policies for bag-of-tasks applications. Besides, for the DP and FSP policies, we also present a method to represent and clusterize functional parameters, like the available amount of FLOPs before a certain deadline. After checking the feasibility of implementing policies with very different objectives on our scheduling model, we have also tested their performance. We have compared them with a scheduler that allocates tasks to nodes at random, and with a centralized version of each policy that has full knowledge of every execution node state. The simpler IBP, MMP and DP policies perform very close to a centralized implementation. On the other hand, the FSP and WDP policies, due to their novelty and the simplifications we have made, present more moderate results. These results open the door to many possible improvements. The immediate one would be to improve the FSP and WDP policies, so that they yield results comparable to the other policies. Then, we could complete the missing parts and create a fully-fledged distributed computing platform, or integrate our code in an existing one. In this way we could test our model on a real scenario. Meanwhile, we want to study the generation and simulation of specific workloads, like many-task, map-reduce and data-intensive applications. Since many policies depend on the task length, we also plan to study different methods of task length estimation. Likewise, availability prediction models could provide better information about execution nodes to those policies where the future state is important, like the MMP and DP policies.
Bibliography 87 [78] RAMAMRITHAM, K., STANKOVIC, J., AND ZHAO, W. Distributed Scheduling of Tasks with Deadlines and Resource Requirements. IEEE Transactions on Computers 38, 8 (Aug. 1989), 1110–1123. [79] RANJAN, R., RAHMAN, M., AND BUYYA, R. A Decentralized and Cooperative Workflow Scheduling Algorithm. In Proceedings of the 8th IEEE International Symposium on Cluster Computing and the Grid (CCGrid) (2008), IEEE, pp. 1–8. [80] RATNASAMY, S., FRANCIS, P., HANDLEY, M., KARP, R., AND SCHENKER, S. A Scalable Content-Addressable Network. In Proceedings of the 2001 Conference on Applications, Technologies, Architectures, and Protocols for Computer Communications (SIGCOMM ’01) (2001), ACM, pp. 161–172. [81] RATNASAMY, S., HANDLEY, M., KARP, R., AND SHENKER, S. Topologically-aware Overlay Construction and Server Selection. In Proceedings of the Twenty-First Annual Joint Conference of the IEEE Computer and Communications Societies (INFOCOM 2002) (2002), vol. 3, IEEE, pp. 1190–1199. [82] RENESSE, R. V., BIRMAN, K. P., AND VOGELS, W. Astrolabe: A Robust and Scalable Technology for Distributed System Monitoring, Management, and Data Mining. ACM Transactions on Computer Systems 21, 2 (2003), 164–206. [83] RHEA, S., GEELS, D., ROSCOE, T., AND KUBIATOWICZ, J. Handling Churn in a DHT. In Proceedings of the USENIX Annual Technical Conference (2004), USENIX Association. [84] RODERO, I., GUIM, F., CORBALAN, J., FONG, L., AND SADJADI, S. M. Grid broker selection strategies using aggregated resource information. Future Generation Computer Systems 26, 1 (2010), 72–86. [85] SCHULZ, S., BLOCHINGER, W., AND HANNAK, H. Capability-Aware Information Aggregation in Peer-to-Peer Grids. Journal of Grid Computing 7, 2 (2009), 135–167. 10.1007/s10723-0089114-z. [86] SHMUELI, E., AND FEITELSON, D. G. On Simulation and Design of Parallel-Systems Schedulers: Are We Doing the Right Thing? IEEE Transactions on Parallel and Distributed Systems 20, 7 (2009), 983–996. [87] SPOONER, D., JARVIS, S., CAO, J., SAINI, S., AND NUDD, G. Local grid scheduling techniques using performance prediction. IEE Proceedings on Computers and Digital Techniques 150, 2 (2003), 87–96. [88] STOICA, I. Chord: a scalable peer-to-peer lookup protocol for internet applications. IEEE/ACM Transactions on Networking 11, 1 (2003), 17. [89] TAKEFUSA, A., CASANOVA, H., MATSUOKA, S., AND BERMAN, F. A Study of Deadline Scheduling for Client Server Systems on Computational Grid. In Proceedings of 10th IEEE International Symposium on High Performance Distributed Computing (2001), IEEE, pp. 406–415.
88 Bibliography [90] TANNENBAUM, T., WRIGHT, D., MILLER, K., AND LIVNY, M. Condor: a distributed job scheduler. In Beowulf cluster computing with Linux. MIT Press, Cambridge, MA, USA, 2002, ch. Condor: a distributed job scheduler, pp. 307–350. [91]The Grid Workloads Archive. http://gwa.ewi.tudelft.nl/pmwiki, nov 2012. [92] VARGA, A. OMNeT++. In Modeling and Tools for Network Simulation. Springer, 2010, pp. 35–59. [93] VISHNUMURTHY, V., CHANDRAKUMAR, S., AND SIRER, E. G. KARMA: A Secure Economic Framework for P2P Resource Sharing. In Proceedings of the Workshop on the Economics of Peerto-Peer Systems (2003). [94] WILHELM, R., ENGBLOM, J., ERMEDAHL, A., HOLSTI, N., THESING, S., WHALLEY, D., BERNAT, G., FERDINAND, C., HECKMANN, R., MITRA, T., MUELLER, F., PUAUT, I., PUSCHNER, P., STASCHULAT, J., AND STENSTRÖM, P. The worst-case execution-time problem–overview of methods and survey of tools. ACM Transactions on Embedded Computing Systems 7, 3 (May 2008), 36:1–36:53. [95] XIAO, L., ZHU, Y., NI, L. M., AND XU, Z. Incentive-Based Scheduling for Market-Like Computational Grids. IEEE Transactions on Parallel and Distributed Systems 19, 7 (2008), 903– 913. [96] YALAGANDULA, P., AND DAHLIN, M. A Scalable Distributed Information Management System. In Proceedings of the Conference on Applications, Technologies, Architectures, and Protocols for Computer Communications (SIGCOMM ’04) (2004), ACM, pp. 379–390. [97] ZHANG, W., AND HE, J. Modeling End-to-End Delay Using Pareto Distribution. In Preceedings of the Second International Conference on Internet Monitoring and Protection (ICIMP) (July 2007), IEEE, p. 21. [98] ZHOU, D., AND LO, V. WaveGrid: a Scalable Fast-turnaround Heterogeneous Peer-based Desktop Grid System. In Proceedings of the 20th IEEE International Parallel and Distributed Processing Symposium (2006), IEEE, pp. 28–37.
Appendix A. Notation Table A.1.: Notation in text and algorithms. Notation Description AiApplication. PRiProperties of application Ai. riRelease time of Ai. eiEnd time of the last task of Ai. niNumber of tasks of Ai. aiLength of a task of Ai, in millions of FLOPs. miRequired memory to execute a task of Ai, in megabytes. diRequired disk space to execute a task of Ai, in megabytes. qiDesired makespan for Ai. δiDeadline of Ai. ziSlowness of Ai. PuPhysical node. Ru/Eu/SuRouting/Execution/Submission node role. AFuAvailability function of node Eu. suComputing power of node Eu, in millions of FLOPs per second. MuMemory available at node Eu, in megabytes. DuDisk available at node Eu, in megabytes. QuQueue end time of node Eu. lu(δ)Amount of FLOPs available before δat Eu. zu(a)Maximum slowness at Euadding one task of length a. x,y,zAvailability summaries. f,g,hSampled functions. f.pParameter of a sampled function. Continues on next page. 89
90 Notation Table A.1.: Continued from previous page. Notation Description SUM(f,g)Sum operation of two sampled functions. DIST(f,g)Distance operation of two sampled functions. τkTask in position kof the queue. βFactor to adjust MMP and FSP forwarding algorithm. νCurrent time. SFmax Maximum summary size. request An application scheduling request. request.srcAddr Request source address. request.PR Application properties of a request. request. pOne of the properties in request.PR. request.nNumber of tasks in a request. Ru.fatherAddr Address of a routing node father. Ru.leftAddr Address of a routing node left child. Ru.rightAddr Address of a routing node right child. Ru.leftInfo Availability information of the left child. Ru.rightInfo Availability information of the right child. Ru.info The aggregation of Ru.leftInfo and Ru.rightInfo. Eu.pA parameter of execution node Eu.
Appendix B. The STaRS Simulator To evaluate the scalability, fault-tolerance and performance of our model, we have developed a simulator in C++. The simulator code is available at http://webdiis.unizar.es/~jcelaya/stars . Our main reason to develop it instead of using an existing alternative was to reduce its memory footprint. In this way, we are able to simulate up to a million nodes with a simple policy, like the IBP. B.1. History The development of a simulator for the STaRS model started in 2005, to obtain the results of [ 32 ]. We built it with Omnet++ [ 92 , 70 ], a DES for computer networks. It includes the INET framework [ 1 ], a good TCP stack and topology model, but we decided to use a simpler network model. The INET framework introduces too much overhead, so we simulated our scheduling model over a star network with fixed link bandwidth and delay. With 2 gigabytes of memory we were able to simulate 50,000 nodes. One of the reasons for using Omnet++ was the possibility of distributing the simulation among several computers, to increase its speed and/or memory usage. However, at that time, it was an experimental feature, it only used a very conservative synchronization protocol and it had no way of correctly stopping the simulation. So, a Master Thesis [15]was carried out to fix these problems. Meanwhile, we started the implementation of a new DES engine. Its objective was to provide only the features of Omnet++ that we needed, stripping everything else out, to reduce the memory footprint as much as possible. Now we are able to simulate one million nodes of the IBP policy in under 10 gigabytes of memory. B.2. Design Our DES engine consists of three main classes: •Simulator: An object of this class drives the simulation. It contains a set of nodes, a network model and an event queue. It iteratively extracts the next event from the queue and sends it to its destination node. The node processes the event, and may generate new events that are inserted in the queue. The network model calculates when the events arrive at its destination. •StarsNode: Each object of this class is a node in the network. It implements the distributed scheduling model presented in this thesis, with the different components of each node role. The model code is decoupled from the simulation engine, so it can be used in a future real 91
92 The STaRS Simulator implementation. The simulator maintains a table with as many instances of this class as nodes in the network. •SimulationCase: An object of this class prepares the simulation and decides how it evolves. For instance, the site-level simulation model presented in Section 8is implemented in a class derived from SimulationCase. At the beginning of the simulation, a SimulationCase instance is created and configured with the contents of a configuration file. In each simulation, a single instance of the Simulator class is created. It then creates a SimulationCase instance and configures it with the contents of a file, that is provided through the command line. This file is just a set of arbitrary key-value pairs. We have a helper tool that generates several configuration files for parameter-sweep simulations. Then, the SimulationCase instance generates the initial events, and the Simulator starts processing the event queue. For the transmission events, the simulator uses a simple network model. It simulates end-to-end links between every pair of nodes, with fixed bandwidth and a random delay that follows a Pareto distribution, as explained in Section 8. It also simulates the transmission and reception queue of every node. The simulation ends when the queue is empty, or certain configurable conditions are met, like reaching a maximum elapsed time or number of submitted requests. Besides these three classes, there are other ones that provide extra functionality. There is an implementation of the different policies in a centralized scheduler, an in-memory implementation of the submission node application database, and several classes that gather data during the simulation, like throughput, traffic and performance statistics. B.3. Future Work Having reduced the memory footprint of the simulator, we are able to simulate huge networks in a reduced amount of memory. However, we have lost important functionality that other simulators provide. The two most important features we miss is a more realistic network model and the distributed simulation. For this reason, we are evaluating the reimplementation of core elements of our DES engine with either Omnet++ or SimGrid [ 30 ]. They provide these features to some extent, and are worth bearing in mind.
Appendix C. Centralized Version of each Policy C.1. IBP Policy The centralized version of the IBP policy is shown in Algorithm C.1. It knows exactly how much memory and disk space is available at every node, and whether they are idle or busy. So, it creates a list of the nodes that fulfill the application requirements and sorts it as in the decentralized version. Then, a task is sent to each of the best nodes as long as there are enough. It is trivial to see that both loops have a time complexity of O ( n ), where n is the number of nodes, and the sort procedure can be performed with complexity O ( nlogn )with a heap structure. Meanwhile, the space complexity is also O(n). C.2. MMP Policy The centralized version of the MMP policy appears in Algorithm C.2. Again, it uses the same heuristic as the decentralized one, allocating tasks to those nodes whose queue will remain shortest afterward. Besides their available memory and disk space, it also records information about the queue end time of every node. The algorithm starts creating a list candidates of nodes that fulfill the application requirements, and their queue end time if they accept a task of the new application. This list is sorted by the queue end time. Then, each task is sent to the node that will finish it earlier. The queue end time of that node is updated, it is inserted in the list and the list is sorted again. The first loop has a time complexity of O ( n ), and the sort procedure can be performed with complexity O ( nlogn )with a heap structure, like in the IBP policy. The second loop has n repetitions at most, but it includes a sorting operation. Since it just inserts a new element in an already sorted heap, its complexity is just O ( logn ). In the end, the whole algorithm can be performed with complexity O(nlogn). The space complexity is still O(n). C.3. DP Policy The centralized version of this policy can be seen in Algorithm C.3. Its workings are similar to the previous two policies. First, it fills the list holes with the nodes that can execute at least one of the new tasks before its deadline. For each node, it calculates the amount of FLOPs that remain free after allocating as many new tasks as possible. Then, the list is sorted so that the nodes with less remaining FLOPs come first, and tasks are allocated to them. The first loop is repeated for every node in the system, but its body contains a call to the function lu ( request.δi ), that returns the amount of FLOPs that can be executed by node Eu before request.δi . 93
94 Centralized Version of each Policy Algorithm C.1 Centralized version of the IBP policy forwarding algorithm. Pre: request is the request, Eis the set of execution nodes. Post: All the tasks in request are allocated to nodes in E. 1: procedure FORWARDCENTRALIZED(request) 2: availableNodes ←; 3: for all Eu∈Edo 4: if Eufulfills request.PR then 5: availableNodes ←availableNodes SEu 6: end if 7: end for 8: SORT(availableNodes).Best nodes are allocated first. 9: while ¬ISEMPTY(availableNodes)Vrequest.n>0do 10: Eu←POPFRONT(availableNodes) 11: task ←EXTRACTTASK(request) 12: SEND(task, Eu) 13: end while 14: end procedure Algorithm C.2 Centralized version of the MMP policy forwarding algorithm. Pre: request is the request, Eis the set of execution nodes. Post: All the tasks in request are allocated to nodes in E. 1: procedure FORWARDCENTRALIZED(request) 2: candidates ←; 3: for all Eu∈Edo 4: if Eufulfills request.PR then 5: candidates ←candidates S{(Eu,Eu.Q+request.ai/su)} 6: end if 7: end for 8: SORT(candidates).Nodes are sorted by their queue end time. 9: while request.n>0do 10: (Eu,te)←POPFRONT(candidates) 11: task ←EXTRACTTASK(request) 12: SEND(task, Eu) 13: PUSHBACK((Eu,te+surequest.ai)) 14: SORT(candidates) 15: end while 16: end procedure
Centralized Version of each Policy 95 Algorithm C.3 Centralized version of the DP policy forwarding algorithm. Pre: request is the request, Eis the set of execution nodes. Post: All the tasks in request are allocated to nodes in E. 1: procedure FORWARDCENTRALIZED(request) 2: holes ←; 3: for all Eu∈Edo 4: if Eufulfills request.PR then 5: h←lu(request.δi) 6: if h≥request.aithen 7: holes ←holes S{(Eu,h,hmod request.ai)} 8: end if 9: end if 10: end for 11: SORT(holes).Nodes are sorted by the remaining amount of FLOPs. 12: while request.n>0do 13: (Eu,th)←POPFRONT(holes) 14: numTasks ←bth/request.ai/suc.Allocate as many tasks as possible 15: tasks ←EXTRACTTASKS(request, numTasks) 16: for all τ∈tasks do 17: SEND(τ,Eu) 18: end for 19: end while 20: end procedure It can be seen in Algorithm 5.1 that its cost is O ( m ), where m is the number of tasks in Eu queue. So, the cost of the first loop is O ( T ), where T is the total number of tasks currently allocated in the system. In practice, this means that the first loop takes much longer to execute than in the two previous policies. Again, sorting the list holes has a time complexity of O ( nlogn ), and the second loop has a complexity of O(n). Finally, in this case, the space complexity is O(T). C.4. FSP Policy Finally, the centralized version of the FSP policy appears in Algorithm C.4. Like the decentralized version, it calculates how many tasks to send to each node to minimize the maximum slowness. However, it works with full knowledge about all the execution nodes. It creates a list candidates with all the nodes that meet the memory and disk requirements, and calculates the slowness obtained by assigning one task to each of them. In this case, the function GETSLOWNESS( Eu,a,n ) uses the algorithms of Section 7.2.1 to calculate the exact value of the maximum slowness reached when adding n tasks of length a to the queue of Eu . Then, it purges the list to keep the request.n best nodes, and iterates on the number of tasks to see if some nodes provide a better slowness with more tasks than others. When no better result is obtained, it sends its assigned number of tasks to each node in the candidates list.
96 Centralized Version of each Policy Algorithm C.4 Centralized version of the FSP policy forwarding algorithm. Pre: request is the request, Eis the set of execution nodes. Post: All the tasks in request are allocated to nodes in E. 1: procedure FORWARDCENTRALIZED(request) 2: candidates ←; 3: for all Eu∈Edo 4: if Eufulfills request.PR then 5: l←GETSLOWNESS(Eu, request.a, 1) 6: candidates ←candidates S{(Eu,l,1)} 7: end if 8: end for 9: PURGE(candidates, request.n) 10: oneMoreTask ←true 11: i←1 12: while oneMoreTask do 13: i←i+1 14: oneMoreTask ←false .If nothing happens, this is the last iteration. 15: for all (Eu,lold,n)∈candidates |n=i−1do 16: lnew ←GETSLOWNESS(Eu, request.a,i) 17: if lnew <candidates.last.slowness then 18: candidates ←candidates \ {(Eu,lold,n)} 19: candidates ←candidates S{(Eu,lnew,i)} 20: oneMoreTask ←true 21: end if 22: end for 23: PURGE(candidates, request.n) 24: end while 25: return candidates.last.slowness 26: end procedure Both loops that iterate on the candidate nodes have a time complexity of O ( n ). The while loop can be repeated up to request.n times, in the case that all the tasks end up in the same node. However, this is extremely rare, and it is repeated no more than four times in most cases. Besides, the candidates list must be kept sorted with a heap structure, to be able to purge and find the worst candidate in O (1). So, like in previous policies, the whole algorithm can be performed with complexity O ( nlogn ). The space complexity is still O(n).