Automatic generation of parallel programs with dynamic load balancing for a network of workstations

Bruce S. Siegell · 1995

Because of their high availability and relatively low cost, networks of workstations are now often considered as platforms for applications that used to be relegated to dedicated multiprocessors. Parallelizing compilers have simplified the programming of shared and distributed memory multiprocessors. However, with networks of workstations, which are more loosely coupled, additional problems of heterogeneity, varying resource availability, and higher communication costs must be addressed in order to maximize utilization of system resources. Computational capabilities may vary with time due to other applications competing for resources, so dynamic load balancing is very important. Our research explores issues in retargeting a parallelizing compiler for a network of workstations. In this dissertation, we describe a system that supports dynamic load balancing of distributed applications consisting of parallelized DOALL and DOACROSS loops. We outline the added compiler functionality needed to generate parallel programs with dynamic load balancing and demonstrate how parameters for dynamic load balancing can be selected and controlled automatically at run time with cooperation between the compiler and runtime system. We have implemented a prototype runtime system on the Nectar system at Carnegie Mellon University and have evaluated its performance using hand-parallelized applications running in various environments. Key performance parameters under our control include the grain size of the application, the frequency of load balancing, and the amount and frequency of work movement. The optimal grain size is selected based on computation and communication costs of the application on the particular system on which it is run. Selecting an appropriate load balancing frequency requires information about communication costs and process scheduling by the operating system. The frequency must be adjusted as loads on the processors change, and controlling the frequency requires the cooperation of the compiler. Making correct decisions regarding work movement is a difficult problem because of high work movement costs and the unpredictable nature of the loads on the processors. Our measurements show that dynamic load performance improves system utilization and reduces execution times in some cases, but is ineffective for others, largely due to the costs of moving work.

Read the paper · More papers on PaperTik