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


Open Problems

  1. Optimally choosing K and queue thresholds is explicitly called an open problem; Tiresias sidesteps it with the foreground–background heuristic (K = 2).
  2. Formal analysis of Discretized 2DAS — precise applicable boundaries in terms of cluster resources and DL job requirements, which would simplify configuration in practice.
  3. 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.
  4. Fine-grained job placement, including intra-server (PCIe) interference, via topology-aware schemes and placement of DL computational graphs.
  5. Dynamic PACKLIMIT determination beyond the current simple linear classifier.