Scalable Load Distribution and Load Balancing for Dynamic Parallel Programs

Emery D. Berger, James C. Browne · 1999

This paper reports design and preliminary evaluation of an integrated load distribution-load balancing algorithm which was targeted to be both efficient and scalable for dynamically structured computations. The computation is represented as a dynamic hierarchical dependence graph. Each node of the graph may be a subgraph or a computation and the number of instances of each node is dynamically determined at runtime. The algorithm combines an initial partitioning of the graph with application of randomized work stealing on the basis of subgraphs to refine imbalances in the initial partitioning and balance the additional computational work generated by runtime instantiation of subgraphs and nodes. Dynamic computations are modeled by an artificial program (k-nary) where the amount of parallelism is parameterized. Experiments on IBM SP2s suggest that the load balancing algorithm is efficient and scalable for parallelism up to 10,000 parallel threads for closely coupled distributed memory architectures.

Read the paper · More papers on PaperTik