Tiresias: A GPU Cluster Manager for Distributed Deep Learning
Juncheng Gu, Mosharaf Chowdhury, Kang G. Shin (Univ. of Michigan, Ann Arbor); Yibo Zhu (Microsoft / Bytedance); Myeongjae Jeon (Microsoft / UNIST); Junjie Qian (Microsoft); Hongqiang Liu (Alibaba); Chuanxiong Guo (Bytedance) | NSDI '19, Feb 26–28 2019, Boston MA | Code: github.com/SymbioticLab/Tiresias
Problem
Distributed deep learning breaks the assumptions of existing cluster managers in three ways: training times are unpredictable, execution is all-or-nothing (every parameter server and worker must be simultaneously active), and GPU sharing is inflexible. Microsoft's production traces show a 10.5× year-over-year increase in DL jobs since 2016, the largest job growing from 32 GPUs in 2016 to 128 GPUs in 2017, while the average queueing delay across all jobs reached 4102 seconds.
Two design failures follow. First, SJF and SRTF need a job's remaining execution time, which predictors like Optimus obtain only by assuming smooth loss curves and run-to-completion — neither holds for AutoML trials that are mostly killed early. Production managers therefore fall back to naïve non-preemptive FIFO; Microsoft's extends Apache YARN's Capacity Scheduler, and users wait up to several hours even for small jobs. Second, managers blindly enforce consolidated placement, packing a job onto the fewest servers with enough GPUs, so a 16-GPU job on 4-GPU machines can block waiting for four completely free servers while GPUs sit idle elsewhere.
Core Insight
A DDL scheduler does not need to know how long a job will run. Priority derived from two-dimensional attained service (GPUs × elapsed time), discretized into a small number of MLFQ-style queues to bound preemption cost, matches perfect-knowledge schedulers on average JCT. Independently, whether a job needs consolidated placement is predicted not by model size but by skew in the model's tensor size distribution, which is externally observable from RDMA-level traffic with no user input.
Method
Scheduling — the 2DAS framework. A job's attained
service is W_J × t_J (GPUs used × running time);
W_J is known at arrival and t_J grows
continuously. The priority function branches on available knowledge:
PRIORITY(Job J, Distribution D):
if D is empty: R_J = -W_J x t_J # LAS, information-agnostic
else: R_J = GittinsIndex(J, D) # partial knowledge
GITTINS_INDEX(Job J, Distribution D):
a_J = W_J x t_J
G_J = sup_{Delta>0} P(S - a_J <= Delta | S > a_J)
-------------------------------------
E[ min{S - a_J, Delta} | S > a_J ]
The Gittins ratio weighs the probability the job completes within service quantum Δ against the expected service it still requires; higher index means higher priority. Under LAS, all jobs start at highest priority and decay as they are served.
Why two dimensions. Single-dimensional policies fail in opposite ways: SRTF lets large jobs with short remaining time hog GPUs and delay small arrivals, while smallest-first lets a stream of small jobs block large jobs near completion. Simulation against shortest-remaining-service-first (SRSF, service = remaining time × GPU count) confirms two dimensions win — normalized w.r.t. SRSF, smallest-first scores 1.52 / 1.20 / 3.45 and SRTF 1.03 / 1.01 / 1.55 on average / median / 95th JCT (>1 is worse).
Priority discretization. Continuous priorities would
preempt constantly, which is expensive on GPUs and degenerates 2DAS into
time-division fair sharing that increases average JCT. Tiresias
keeps K logical queues Q1...QK, where
queue i holds jobs with W_J t_J in
[Q_i^lo, Q_i^hi), Q_1^lo = 0,
Q_K^hi = ∞. New jobs enter Q1 and are demoted
when attained service crosses Q_i^hi. Under LAS, same-queue
jobs run FIFO by start time — not submission time,
because all-or-nothing allocation forces GPU-starved high-priority jobs
to be skipped over. Under Gittins, Δ_i = Q_i^hi and
same-queue jobs sort by index; in the last queue Δ_K = ∞,
so Gittins degenerates to LAS. Following the classic
foreground–background result for heavy-tailed distributions, Tiresias
uses K = 2.
Starvation avoidance. A job waiting longer than
STARVELIMIT is promoted to Q1. One operator
knob controls the trade-off, with t_J and δ_J
reset on promotion so the job is not immediately demoted:
δ_J >= PROMOTEKNOB * t_J
PROMOTEKNOB = ∞ disables promotion and purely minimizes
average JCT; smaller values buy fairness at the cost of average JCT.
Placement — profile-based, not ILP. An ILP
minimizing and balancing inter-machine traffic was built first and
rejected: too slow at cluster scale, and explicitly balancing network
load does not necessarily improve training performance. Instead the
profiler intercepts RDMA ibverbs APIs on every server
(no RDMA-level traffic monitor existed, so they built one) to measure
the skew S_J of tensors across parameter servers,
framework-agnostically. Because every iteration is
communication-identical, a few iterations of profiling suffice. If
S_J > PACKLIMIT, consolidate; otherwise spread to reduce
fragmentation. PACKLIMIT is refreshed periodically by a simple
linear classifier over job placement and observed
performance.
The mechanism behind skew: each TensorFlow tensor is wrapped as one communication message, so message-size distribution follows tensor-size distribution. A few huge tensors suffer network contention badly; many small tensors interleave well. Concretely, the VGG family and AlexNet are dominated by a single tensor (61.0%–77.3% of model size), whereas ResNet, Inception, and GoogleNet are far less skewed (3.6%–14.6%) — exactly the split observed in the placement-sensitivity measurements.
Experimental Setup
| Component | Value |
|---|---|
| Machines | 15 × IBM PowerNV 8335-GTB, Michigan ConFlux cluster |
| GPUs | 4 × NVIDIA Tesla P100 (16 GB) per machine; 60 total |
| CPU / memory | 2 × 10-core POWER8 (8 threads/core); 256 GB DDR4 |
| Network | 100 Gbps EDR Mellanox InfiniBand |
| Shared FS | GPFS, 1.2 GB/s read and write per machine; holds checkpoints |
| Framework | TensorFlow 1.3.1 with RDMA extension, TF benchmark models |
| Workload | 480 DL/DDL jobs scaled down from the Microsoft trace |
| Job mix | half single-GPU; 40 2-GPU, 80 4-GPU, 90 8-GPU, 25 16-GPU, 5 32-GPU |
| Parameter servers | equal to each job's GPU count |
| Training time | 2 minutes to 2 hours; fixed iteration count per job |
| Arrivals | Poisson, average inter-arrival 30 seconds |
| Tiresias config | K = 2, threshold 3200 GPU seconds, PROMOTEKNOB disabled |
| Baselines | YARN-CS, SRTF (complete information + Tiresias placement), Gandiva, Best-effort |
| Metric | Factor of Improvement = Duration of an Approach / Duration of Tiresias-L |
Job bins by GPU count and training time (small ≤ 4 GPUs, short < 800 s after scale-down): Bin 1 (SS) 63.5%, Bin 2 (SL) 12.5%, Bin 3 (LS) 16.5%, Bin 4 (LL) 7.5%. A discrete-time simulator replays the full Microsoft trace at scale, using actual job completion times because it cannot model training time under dynamic cluster conditions.
Headline Quantitative Results
Tiresias-L = Discretized 2D-LAS; Tiresias-G = Discretized 2D-Gittins index.
Testbed (60 GPUs). Average JCT is 5.5× better than YARN-CS; median JCT 27× better. Tiresias-G is essentially identical to Tiresias-L (1.06× average, 1.05× median), its small loss caused by more preemptions. Bin 1 (small-short) average JCT is 300 s (Tiresias-L) and 330 s (Tiresias-G) — 27.6× and 25.2× better than YARN-CS — while Bin 4 (large-long) jobs see almost the same average JCT under Tiresias and YARN-CS. GPU-utilization CDFs are similar across all solutions, but makespan differs: Tiresias-L 27,400 s, Tiresias-G 27,510 s, SRTF 28,070 s, YARN-CS 33,270 s — a 1.21× makespan reduction over YARN-CS.
Queueing delays for DDL jobs:
| Solution | Average | Median | 95th |
|---|---|---|---|
| YARN-CS | 8146 s | 7464 s | 15327 s |
| SRTF | 593 s | 32 s | 3133 s |
| Tiresias-G | 1005 s | 39 s | 7933 s |
| Tiresias-L | 963 s | 13 s | 7755 s |
Preemption overhead: Tiresias-L 221 preemptions / 13,724 s; Tiresias-G 297 / 17,425 s; SRTF 316 / 18,057 s. Tiresias-G preempts more because same-queue jobs are re-sorted by Gittins index at every event, while Tiresias-L never re-sorts under FIFO.
Placement contribution: rerunning Tiresias-L without the profiler (random placement) shows up to 1.67× improvement from profile-based placement, with fewer than 30% of DDL jobs experiencing limited performance loss. Scheduling — not placement — remains the dominant source of gain.
Trace-driven simulation of the Microsoft trace (normalized by Tiresias-L):
| Solution | Average | Median | 95th |
|---|---|---|---|
| YARN-CS | 2.41× | 30.85× | 1.25× |
| Best-effort | 1.50× | 9.03× | 1.08× |
| SRTF | 1.00× | 1.00× | 0.84× |
| Gandiva | 2.00× | 2.59× | 2.08× |
| Tiresias-G | 0.97× | 1.00× | 0.85× |
Tiresias-L beats YARN-CS by 2.4×, Best-effort by 1.5×, and Gandiva by 2× on average JCT, and matches SRTF, which uses complete knowledge — SRTF is better only at the 95th percentile (0.84×). Simulator fidelity check on the testbed workload: 5.11× average / 1.50× 95th w.r.t. YARN-CS.
Sensitivity. Thresholds ≥ 1 h GPU time are flat for Tiresias-L (0.5 h is 1.03×), because 1 hour of GPU time covers more than 60% of jobs; K = 3 or 4 improves average JCT by only 1% over K = 2; finite PROMOTEKNOB cuts Tiresias-G's maximum JCT to 0.65× but moves Tiresias-L's 95th JCT only from 1.000 to 0.936, the trace being heavy-tailed.
Preemption cost tracks tensor count, not bytes: VGG19 (548.1 MB, 39 tensors) checkpoints in 2.7 s, while ResNet152 (230.2 MB, 778 tensors) takes 26.3 s.
Limitations
- No formal analysis of Discretized 2DAS's applicable boundaries; K, queue thresholds, PACKLIMIT, and PROMOTEKNOB are all set empirically.
- Preemption remains checkpoint-based and expensive — 221–297 preemptions cost 13,724–17,425 s of wall-clock overhead in a single testbed run.
- Placement is coarse-grained (consolidate vs. distribute only) and ignores intra-server interference such as PCIe contention among collocated workers and parameter servers.
- Parameter-server architecture only; all-reduce aggregation is not evaluated, and the skew argument is framed entirely around tensor-to-PS mapping.
- Data parallelism and synchronous training only; model parallelism is explicitly out of scope.
- The simulator cannot capture preemption overhead, placement impact, or cluster dynamics and replays actual completion times instead — hence the gap between testbed (5.5×) and simulated (5.11×) numbers.
- Scale gap: 60 P100 GPUs on 15 machines versus production Philly's 100 4-GPU plus 250+ 8-GPU servers, with a scaled-down 480-job workload rather than the trace itself.
- The Gittins index needs a duration distribution that stays valid in the near future; no drift-detection mechanism is given.
Open Problems
- Optimally choosing K and queue thresholds is explicitly called an open problem; Tiresias sidesteps it with the foreground–background heuristic (K = 2).
- Formal analysis of Discretized 2DAS — precise applicable boundaries in terms of cluster resources and DL job requirements, which would simplify configuration in practice.
- Lightweight preemption. Existing DL preemption primitives are time-consuming; Gandiva's alternative requires framework modification and still carries non-negligible overhead. With cheap preemption, many classic network-flow and CPU-scheduling algorithms become applicable to DDL.
- Fine-grained job placement, including intra-server (PCIe) interference, via topology-aware schemes and placement of DL computational graphs.
- Dynamic PACKLIMIT determination beyond the current simple linear classifier.