Scheduling parallel programs in distributed systems
C. Gary Rommel, John A. Stankovic · 1988
This thesis investigates scheduling parallel programs under the processor-sharing discipline for uniprocessors, for multiprocessors, and for distributed systems. Two classes of parallel programs are considered: those without any IPC (called Fork-Join jobs) and those with asynchronous and uniform IPC (called clusters). We divide our study into two parts. The first develops analytical solutions for Fork-Join Jobs on uniprocessors and multiprocessors; and the second develops and evaluates via simulation Fork-Join jobs and clusters on distributed systems. In the first part of the thesis, the types of site scheduling studied are TS-PS where tasks of a job are scheduled independently at processor-sharing servers, JS-PS in which tasks of a job are scheduled as a single entity at processor-sharing servers, and FCFS where tasks of a job are scheduled independently by order of arrival. For Poisson job arrivals and exponentially distributed task service times, we found analytical solutions and computationally efficient bounds for Fork-Join TS-PS and JS-PS job response times. We observed that scheduling by FCFS is better than TS-PS or JS-PS, scheduling by JS-PS is better than TS-PS on uniprocessors, but for multiprocessors scheduling by TS-PS is better than scheduling by JS-PS under moderate site utilizations. Partitioning the multiprocessors of a system into two disjoint sets, one for single task jobs and one for Fork-Join jobs, was found never to improve the response time of both classes of jobs, but in some cases was found to degrade the response time of both classes of jobs. We developed an algorithm to schedule parallel programs on distributed systems. Over a wide range of parameters our algorithm was found to be superior to both no load balancing, NLB, and shortest queue first scheduling, SQF. Through performance studies our algorithm was determined to be much better than both NLB and SQF scheduling as the intra-job communication increased, as the number of tasks per job increased, as the variance of tasks per job increased, and as the site utilization increased. Our algorithm was found to work well in the presence of imperfect information, and was found to work well in task transfers across the network. (Abstract shortened with permission of author.)