Market-Like Task Scheduling in Distributed Computing Environments

Thomas W. Malone · DSpace@MIT (Massachusetts Institute of Technology) · 2011

This paper focuses on a class of market-like methods for decentralized scheduling of tasks in distributed computing networks.In these methods, processors send out "requests for bids" on tasks to be done and other processors respond with "bids" giving estimated completion times that can reflect such factors as machine speed and data availability.A simple and general protocol for such scheduling is described and simulated in a wide variety of situations (e.g., different network configurations, system loads, and message delay times).The protocol is found to provide substantial performance improvements over processing tasks on the machines at which they originate even in the face of relatively large message delays and relatively innaccurate estimates of processing times.The protocol also performs well in comparison to one simpler and one more complex alternative.In the final section of the paper, a prototype system is described that uses this protocol for sharing tasks among personal workstations on a local area network.Market-like Task SchedulinginDistributedComputingEnvironmentsWith the rapid spread of personal computer networks and the increasing availability of low cost VLSI processors, the opportunities for massive use of parallel and distributed computing are becoming more and more compelling ([JonSO); [Gaj85]; [Dav 81bl; [Ens8l]; [Ber82I; [Bir821).One of the fundamental problems that must be solved by all such systems is the problem of how to schedule tasks on processors.This problem is, of course, a well-known one in traditional operating systems and scheduling theory, and there are a number of mathematical and software techniques for solving it in both single and multi-processor systems (e.g.,[Bri73]; [Cof73]; [Con671; [JonSO]; [KriTlJ; [Lam68);[Kle81];[Wit80];[Vnt811). Almost all the traditional work in this area, however, deals with centralized scheduling techniques, where all the information is brought to one place and the decisions are made there.In highly parallel systems where the information used in scheduling and the resulting actions are distributed over a number of different processors, there may be substantial benefits from developing decentralized scheduling techniques.For example, when a centralized scheduler fails, the entire system is brought to a halt, but systems that use decentralized scheduling techniques can continue to operate with all the remaining nodes.Furthermore, much of the information used in scheduling is inherently distributed and rapidly changing (eg , momentary system load).Thus, decentralized scheduling techniques can "bring the decisions to the information" rather than having to constantly transmit the information to a centralized decision maker.Because of these advantages, a growing body of recent work has begun to explore such decentralized scheduling techniques in more detail (e.g., [StaSSJ;[Sin85]; [Liv82]; [Mal83|; [Lar821; [SmiSO], [Sto771; [Sto781; [Cho79]; [Bry81I; [Sta84]; [Ten81al, [TenSlb]; [StrSl] ; [SRC851; see Stankovic for a useful review).In this paper, we focus on a particular class of decentralized scheduling techniques: those involving market-like "bidding" mechanisms to assign tasks to processors (eg, [Mal831;[Far72];[SmiSO]; [Sin85]; [Far73]; [Sta84]).Such techniques are remarkably flexible in terms of the kinds of factorsthey can take into account: job characteristics, processor capacities and speeds, current network loading and current locations of data and related tasks (e.g., see [Sta84]).In succeeding sections of the paper, we will describe a simple but powerful bidding protocol and present detailed simulation and analytic results that explore the behavior of the protocol in a wide variety of stituations.These results apply to many forms of parallel computation, regardless of whether or not the processors are geographically separated and whether or not they share memory. Motivating exampleEven though the simulation results are applicable in many situations, the driving example that motivated the development and analysis of our protocol was the increasingly common situation of large numbers of personal workstations connected by local area networks (e.g., [Bir82I; [BogSO]).In the final section of the paper, we describe a prototype system, called Enterprise, that uses the protocol to share tasks among workstations in such a network.One of the benefits of sharing tasks in such networks is that a new philosophy for designing distributed systems becomes possible.The traditional philosophy used in designing distributed systems based on local area networks is to have dedicated personal workstations which remain idle when not used by their owners, and dedicated special purpose servers such as file servers, print servers, and various kinds of data base servers (e.g., [Bir821;[Sch84]).A system like Enterprise that schedules tasks on the best processor available at run time (either remote or local) enables a much more flexible design.In this new philosophy, personal workstations are still dedicated to their owners, but during the (often substantial) periods of the day when their owners are not using them, these personal workstations become general purpose servers, available to other users on the network."Server" functions can migrate and replicate as needed on otherwise unused machines (except for those such as file servers and print servers that are required to run on specific machines).Thus programs can be written to take advantage of the maximum amount of processing power and parallelism available on a network at any time, with little extra cost when there are few extra machines available. Problem description and related workIn describing the problem on which we are focusing, it is useful, first of all, to distinguish between two components of task scheduling in distributed networks: (1) the assignment of tasks to processors, and(2) the sequencing of task once they have been assigned to processors.Much previous work on distributed load sharing has focused only on the task assignment problem (eg, [Liv82); [Sta85|;[ Cho79]) and used first come-first served (FCFS) sequencing.It is clear, however, that task sequencing can have a major effect on system performance measures (e.g., see [Con67|).The task scheduling method we will describe solves both the task assignment and the task sequencing problems simultaneously.

Read the paper · More papers on PaperTik