Automatic partitioning and scheduling on a network of personal computers

David Andrew Hornig · 1984

This dissertation investigates automatic partitioning, scheduling, and fault tolerance for distributed applications on a network of personal computers. Previous distribution systems in loosely coupled environments have used a special process control language or special language features to describe the task structure of a program while the lower level algorithms were defined in a standard algorithmic language. Like many of the systems proposed for tightly coupled processors, we use a single applicative language for both the task structure and the algorithm, doing away with the distinction between tasks and functions by making every function a potential task. However, the network environment has deeply influenced the design: the long communication latencies commit the system to large grain parallelism and carefully planned communication strategies. The product of the research is a language and execution environment called STARDUST. STARDUST programs are written in a general-purpose applicative language and marked with an estimate of each function's execution time. The system partitions the programs by expanding user function calls with high execution time estimates and by breaking up calls to the system's list operators. The segments are scheduled on the available processors, taking into account both load balancing and the costs of message passing. Failed tasks can be restarted by moving them to surviving processors. Three sets of experiments run on Perq computers and an Ethernet gave two successes and one failure. The quick sort experiment failed to achieve significant speedups due to some poor scheduling heuristics and large amounts of overhead that could not be transferred to other processors. The six-processor versions of the signal processing and molecular modelling experiments showed speedups of about two when the exported function calls took one second to evaluate, and speedups of over three when the exported calls took five seconds. Message passing and interpreter overhead are the main reasons that such a large granularity is needed; partitioning and scheduling do not contribute significantly to overall execution time. The experimental results also include a demonstration of automatic failure recovery and run-time redundant subcomputation elimination.

Read the paper · More papers on PaperTik