Tiresias: A GPU Cluster Manager for Distributed Deep Learning — Detailed Summary

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 | ISBN 978-1-931971-49-2 | Code: github.com/SymbioticLab/Tiresias

Per-section summary following the paper's own headings. Every paragraph produces at least one bullet; all tables, equations, and figure values are preserved exactly.


Abstract


1. Introduction


2. Background and Motivation

2.1 Distributed Deep Learning

2.2 Challenges

Figure 4 — checkpoint (pause) time in TensorFlow, chief worker only:

Model VGG19 VGG16 VGG11 AlexNet ResNet152 ResNet101 ResNet50 Inception4 Inception3 GoogleNet
Checkpoint time (s) 2.7 2.3 1.9 1.3 26.3 17.8 9.4 22.2 14.1 5.7

2.3 Potential for Benefits


3. Tiresias Design

3.1 Overall Architecture

3.2 Scheduling

3.2.1 Why Two-Dimensional Scheduling?

Table 1 — Normalized performance of single-dimensional schedulers w.r.t. SRSF:

Scheduler Avg. JCT Med. JCT 95th JCT
Smallest-First (SF) 1.52 1.20 3.45
SRTF 1.03 1.01 1.55

3.2.2 Two-Dimensional Attained Service Scheduler (2DAS)

Pseudocode 1 — Priority function in 2DAS:

1: procedure PRIORITY(Job J, Distribution D)
2:     if D is empty then            # w/o distribution, apply LAS
3:         R_J = -W_J x t_J
4:     else                          # w/ distribution, apply Gittins index
5:         R_J = GittinsIndex(J, D)
6:     return R_J

9:  procedure GITTINS_INDEX(Job J, Distribution D)
10:     a_J = W_J x t_J
11:     G_J = sup_{Delta>0}   P(S - a_J <= Delta | S > a_J)
                             -------------------------------------
                             E[ min{S - a_J, Delta} | S > a_J ]
12:     # P = probability, E = mean, both from D; Delta = service quantum
13:     return G_J

3.2.3 Priority Discretization

δ_J  >=  PROMOTEKNOB * t_J

3.3 Placement

3.3.1 Profiler

3.3.2 The Placement Algorithm

3.4 Summary

Table 2 — Comparison of DL cluster managers:

Dimension YARN-CS Gandiva Optimus Tiresias (Gittins index) Tiresias (LAS)
Prior Knowledge None None JCT prediction JCT distribution None
Scheduling Algorithm FIFO Time-sharing Remaining-time-driven Gittins index LAS
Scheduling Input Arrival time N/A Remaining time Attained service Attained service
Schedule Dimensions Temporal None Temporal Spatial & temporal Spatial & temporal
Job Priority Continuous Continuous Continuous Discretized queues Discretized queues
Job Preemption N/A Context switch Model checkpoint Model checkpoint Model checkpoint
Minimizing Average JCT No No Yes Yes Yes
Starvation Avoidance N/A N/A Dynamic resource Promote to Q1 Promote to Q1
Job Placement Consolidation Trial-and-error Capacity-based Profile-based Profile-based

Pseudocode 2 — 2DAS Scheduler:

1:  procedure 2D-LAS(Jobs J, Queues Q1...QK, Distribution D)
2:      P = {}                                # jobs to preempt
3:      for all Job J in J do
4:          if J is RUNNING then
5:              r_J = PRIORITY(J, D)
6:          if J is WAITING longer than STARVELIMIT then
7:              Reset t_J
8:              Enqueue J to Q1               # promote if STARVING
9:      while Cluster has available GPUs do
10:         for all i in [1,K] do             # prioritize across queues
11:             if D is not empty and i in [1, K-1] then
12:                 Sort GittinsIndex(Qi)
13:             for all Job J in Qi do
14:                 if Available GPUs >= W_J then
15:                     Mark W_J GPUs as unavailable
16:                 else
17:                     P = P union J
18:                     Preempt J if already RUNNING
19:     for all Job J in J and J not in P do
20:         if J is not already RUNNING then
21:             if J was not profiled before then
22:                 Profile J                          # Sec 3.3.1
23:             Store J's start time     # FIFO key for Discretized 2D-LAS
24:             Assign GPUs by comparing S_J to PACKLIMIT    # Sec 3.3

4. Implementation


5. Evaluation

Stated highlights: on the testbed Tiresias improves average JCT by up to 5.5× and makespan by 1.21× vs. YARN-CS and performs comparably to SRTF with complete prior information (§5.2), driven by placement benefits for skewed jobs and reduced queueing delay; benefits hold in large-scale simulation of the Microsoft trace (§5.3); and Tiresias is robust to configuration parameters and workload variation (§5.4). Tiresias-G = Discretized 2D-Gittins index; Tiresias-L = Discretized 2D-LAS.

5.1 Experimental Setup

Component Value
Machines 15 × IBM PowerNV 8335-GTB, Michigan ConFlux cluster
GPUs per machine 4 × NVIDIA Tesla P100, 16 GB GPU memory
Total GPUs 60
CPU 2 × 10-core POWER8 (8 threads per core)
Memory 256 GB DDR4
Network 100 Gbps EDR Mellanox InfiniBand adapter
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
Per-model jobs 48 jobs per model in Table 6
Training time 2 minutes to 2 hours; fixed iteration count per job
Arrivals Poisson, average inter-arrival 30 seconds
Metric Factor of Improvement = Duration of an Approach / Duration of Tiresias-L

Table 3 — Job bins by spatial (GPU count) and temporal (training time):

Bin 1 (SS) 2 (SL) 3 (LS) 4 (LL)
% of Jobs 63.5% 12.5% 16.5% 7.5%

5.2 Testbed Experiments

5.2.1 JCT Improvements

5.2.2 Cluster-Wide GPU Utilization

Scheduler Makespan
Tiresias-L 27,400 s
Tiresias-G 27,510 s
SRTF 28,070 s
YARN-CS 33,270 s

5.2.3 Sources of Improvements

Table 4 — 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

5.2.4 Overheads

Scheduler Preemptions Total preemption time
Tiresias-L 221 13,724 s
Tiresias-G 297 17,425 s
SRTF 316 18,057 s

5.3 Trace-Driven Simulations

Table 5 — Improvements in JCT in simulation (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×

5.4 Sensitivity Analysis

Notation (K, threshold1, threshold2, ...); e.g. (2, 1h) = two queues with a 1-hour-GPU-time threshold.

Queue thresholds (§5.4.1) — Figures 14a / 15a, normalized average JCT w.r.t. (2, 1h):

Threshold / Δ 0.5h 1h 2h 4h
Tiresias-L 1.03 1.00 1.00 1.00
Tiresias-G 1.00 1.00 1.00 1.00

Number of queues K (§5.4.2) — Figures 14b / 15b, with settings (2, 1h), (3, 1h, 2h), (4, 1h, 2h, 4h):

K 2 3 4
Tiresias-L 1.00 0.99 0.99
Tiresias-G 1.00 1.00 1.00

PROMOTEKNOB (§5.4.3) — based on (2, 1h), starting at 1 and increasing by powers of 2; disables promotion, smaller values promote long-delayed jobs more often:

PROMOTEKNOB Inf 8 4 2 1
Tiresias-L, norm. 95th JCT 1.000 0.952 0.947 0.943 0.936
Tiresias-G, norm. max JCT 1.00 0.65 0.65 0.65 0.65

6. Discussion and Future Work



8. Conclusion


Appendix A — Characteristics of the Production Cluster


Table 6 — Characteristics of 10 popular DNN models in TensorFlow:

Model Model size (MB) #Tensors #Large tensors (≥1MB) Largest tensor size (MB) Largest tensor ratio
VGG19 548.1 39 15 392.0 71.5%
VGG16 527.8 33 12 392.0 74.3%
VGG11 506.8 23 9 392.0 77.3%
AlexNet 235.9 17 7 144.0 61.0%
ResNet152 230.2 778 48 9.0 3.9%
ResNet101 170.4 523 35 9.0 5.3%
ResNet50 97.7 268 18 9.0 9.2%
Inception4 162.9 599 81 5.9 3.6%
Inception3 91.0 397 21 7.8 8.6%
GoogleNet 26.7 117 7 3.9 14.6%

Appendix C — 2D-Gittins Index Value


Appendix D — ILP Formulation for DDL Placement

Symbol Meaning
N number of GPU nodes in the cluster
N_i the i-th node
t_i existing network traffic on N_i from current DDL jobs
g_i free GPUs on N_i before placing any new job
J, M new DDL job to place, and its model size
W, K number of GPU workers and parameter servers in J
s_j total size of tensors hosted by the j-th parameter server
w_i number of J's workers placed on N_i
p_ji binary: 1 if the j-th PS is placed on N_i, else 0

Total traffic on N_i = existing traffic + traffic from J's workers on it + traffic from J's parameter servers on it, minus traffic between collocated PS and workers:

T_i = t_i + w_i * ( M - SUM_{j in K} p_ji * s_j )
          + SUM_{j in K} p_ji * s_j * ( W - w_i )

minimize  max_{i in N} { T_i }

(1)  for all i in N:   w_i <= g_i             # GPU resource limit per node
(2)  SUM_{i in N} w_i = W                     # all workers of J are placed
(3)  for all j in K:   SUM_{i in N} p_ji = 1  # each PS has exactly one host

Limitations (stated or implied by the paper)