Online Scheduling of Heterogeneous Distributed ML Jobs: Beyond Static Allocation

Online scheduling of heterogeneous distributed machine learning jobs

2020-10-08
Qin Zhang, Ruiting Zhou, Chuan Wu, Lei Jiao, Zongpeng Li
Summary
Problem
Method
Results
Takeaways
Abstract

This paper presents a novel online scheduling algorithm for heterogeneous distributed machine learning (ML) jobs in Parameter Server (PS) architectures. The core method utilizes demand elasticity to dynamically determine execution windows and resource configurations, achieving state-of-the-art performance in minimizing weighted average completion time.

TL;DR

Researchers have developed a sophisticated online scheduling algorithm that treats Machine Learning training as an elastic process. Unlike previous schedulers that assume more GPUs always mean linear speedup, this approach understands the law of diminishing returns in parallel training. By using a primal-dual framework and an interval-based scheduling strategy, it reduces job completion time by up to 200% compared to standard industry methods.

Background: The Hidden Complexity of ML Scheduling

In modern AI clouds, managing distributed ML jobs is a balancing act. Most production systems like Google’s Borg or Microsoft’s YARN variants rely on users to specify their resource needs. However, ML jobs are inherently resource-elastic. For instance, doubling the GPUs from 1 to 2 might reduce training time from 15ms per batch to 10ms (not 7.5ms) due to the overhead of synchronizing gradients.

The fundamental challenge addressed here is: How can a cluster operator decide the exact number of Workers and Parameter Servers (PS) for an arriving job to maximize total cluster efficiency?

The "Elasticity" Insight

The author's core insight is that by intelligently "under-provisioning" or "right-sizing" jobs based on their training efficiency curves, the total system-wide average completion time can be significantly lower.

If we have 3 GPUs and two identical jobs arrive, should we give all 3 to Job A and then 3 to Job B? Or give 1 to A and 2 to B simultaneously? The latter often wins because the marginal benefit of the 3rd GPU on a single job is often less than its benefit as the 1st GPU for a new job.

Methodology: Primal-Dual and Time-Intervals

The solution is structured into two elegant layers:

1. Online Framework (Interval-based)

Since we don't know the future, the algorithm partitions the total timespan into intervals (). It groups jobs arriving in one interval into a "batch" and schedules them as a sub-problem. This converts a continuous online challenge into a series of reachable offline optimization goals.

2. Primal-Dual Batch Scheduling

To solve the batch problem, the authors use a Primal-Dual architecture.

  • Dual Variables as Prices: Every resource (CPU, GPU, Bandwidth) on every server is assigned a virtual "price."
  • Exponential Cost Functions: As a server gets busier, the price of its resources increases exponentially.
  • Utility Maximization: A job is only scheduled if its "Priority Weight" outweighs the "Resource Cost" of its best possible configuration.

Overall Architecture Figure: The formula for processing capacity accounts for both computation and communication overhead .

Experimental Evidence: SOTA Performance

The algorithm was tested against FIFO, DRF (Dominant Resource Fairness), and OASiS.

  • Weighted Completion Time: The proposed consistently achieved lower completion times, especially as the system became more saturated (resource shortage).
  • Scalability: The computational complexity remains polynomial, making it practical for real-time cloud management.

Performance Comparison Figure: Total weighted completion time comparison across different job counts.

Critical Analysis & Takeaways

The brilliance of this work lies in its mathematical rigor—it doesn't just provide a "heuristic"; it provides an algorithm with a bounded competitive ratio. This means we can mathematically guarantee that the online performance won't be catastrophically worse than a hypothetical "God-mode" (offline) scheduler that knows the future.

Limitations:

  • Synchronous Training Focus: The model is optimized for S-SGD. While S-SGD is dominant, the rise of Asynchronous and Decentralized (All-Reduce) training might require adjustments to the communication cost model.
  • No Preemption: The current model assumes jobs run to completion without interruption. In real-world clouds, preemption (suspending a job to run a higher priority one) is common, though it introduces significant Checkpoint/Restart overhead.

Future Outlook

This research highlights the necessity for "Application-Aware" scheduling. The next generation of AI clouds won't just look at "Available CPUs"; they will look at the convergence rate of the model being trained to make the most efficient allocation decisions.

Find Similar Papers

Try Our Examples

  • Search for recent papers that extend online ML job scheduling by incorporating job preemption and migration to further reduce resource fragmentation.
  • What are the key differences in resource elasticity modeling between synchronous SGD (this paper) and asynchronous training methods in distributed ML clusters?
  • Explore newer literature that applies primal-dual online optimization frameworks to All-Reduce based distributed training architectures instead of the Parameter Server model.
Contents
Online Scheduling of Heterogeneous Distributed ML Jobs: Beyond Static Allocation
1. TL;DR
2. Background: The Hidden Complexity of ML Scheduling
3. The "Elasticity" Insight
4. Methodology: Primal-Dual and Time-Intervals
4.1. 1. Online Framework (Interval-based)
4.2. 2. Primal-Dual Batch Scheduling
5. Experimental Evidence: SOTA Performance
6. Critical Analysis & Takeaways
7. Future Outlook