Online Scheduling of Heterogeneous Distributed ML Jobs: Beyond Static Allocation
Online scheduling of heterogeneous distributed machine learning jobs
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.
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.
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.
