Skip to main content

HiveD: Sharing a GPU Cluster for Deep Learning with Guarantees

HiveD reserves groups of GPUs that sit together, not GPU counts, so a tenant never queues longer in a shared cluster than in a private one of the same size.

TL;DR

  • In a shared GPU cluster a tenant reserves a quota, a number of GPUs. Deep learning jobs need more than a number: an 8-GPU job wants its GPUs on one node. When other tenants' small jobs are spread over every node, a tenant can be inside its quota and still unable to start a job that would start at once on a private cluster of the same size. The paper calls this the sharing anomaly.
  • HiveD replaces the quota with cells: reserved groups of GPUs at a level of the hardware hierarchy (one GPU, the GPUs under one PCIe switch, under one CPU socket, in one node, in one rack). A tenant's cells form a virtual private cluster, and any existing scheduler can run inside it.
  • Cells of the virtual clusters are bound to physical GPUs on demand by buddy cell allocation, which splits a larger cell when a smaller one is needed and merges cells back when they are released. The paper proves that every request within a tenant's reservation can always be satisfied.
  • With three existing schedulers on a 96-GPU cluster, the extra queuing delay of one tenant falls from as much as 1,000 minutes to zero, at about the same job completion time. In simulation on a two-month production trace, a tenant that queues nearly 7ร— longer under quota than in its private cluster queues less than in its private cluster under HiveD.

The sharing anomaly

Large organizations run one GPU cluster for many teams. Each team, a tenant, contributes budget or hardware and receives a quota in return: the right to use that many GPUs, with the CPU and memory that go with them. The understanding is that a tenant can always get at least the share it paid for (Section 2).

A training job asks for more than a count. Its GPUs should sit close together, because a job on eight GPUs in one node usually runs much faster than the same job on eight nodes (the hardware behind that is on the multi-GPU communication page). The paper writes the shape of a request as nodes ร— GPUs per node: a 64-GPU job wants 8ร—8, eight nodes with eight GPUs each, and not 64ร—1. A requirement is hard if the job waits until that shape is free, and soft if it accepts a worse placement and runs slower.

The observation that started the work came from user complaints in a production cluster. A tenant with a quota of 64 GPUs could not run a single 8ร—8 job, the only job it had. The GPUs that would have given it that shape had been fragmented by other tenants' jobs. The tenant had the quota, and its job still had to wait or run in a worse shape.

The figure replays the complaint at a small scale: two tenants with a quota of 8 GPUs each, on two 8-GPU nodes. Step through it under quota, then switch to cells.

The paper compares this to external fragmentation in memory management, with one difference that matters. A program suffers from its own fragmentation, but here the fragmentation is caused by other tenants, and the tenant that suffers can do nothing about it. Quota reserves a quantity of GPUs and says nothing about where they are.

The definition follows from the complaint. A cluster suffers from sharing anomaly if a tenant's sequence of requests cannot be satisfied in the shared cluster and could be satisfied in a private cluster as large as the tenant's quota. A cluster that never does this has sharing safety. The paper reports that important users went back to private clusters after long queuing delays of this kind, and that tenants with large reservations suffer the most.

One answer would be a scheduler that keeps global fragmentation low. The authors argue against it: schedulers for deep learning already balance several goals, and packing jobs tightly can slow them down through interference. HiveD takes the guarantee out of the scheduler. It is a reservation layer that enforces sharing safety, and it leaves utilization, job completion time and fairness to whatever scheduler runs on top.

Cells: reserving a shape

A cell is a group of GPUs at one level of the hardware hierarchy, together with their interconnect. In the paper's example there are four levels: a single GPU (level 1), the two GPUs under one PCIe switch (level 2), the four GPUs under one CPU socket (level 3) and the eight GPUs of a node (level 4). A cell of one level consists of buddy cells of the level below, and buddies can be merged back into the cell above them.

A tenant reserves a number of cells at each level. That reservation is its virtual private cluster (VC). The paper's example is one rack of four 8-GPU nodes shared by three tenants (Figure 3). Pick a tenant to see what it reserved and what its scheduler sees.

To the scheduler inside a VC, the cells are a private cluster made of nodes of different sizes. Tenant C's scheduler sees two 8-GPU nodes and one 2-GPU node, and need not know that the 2-GPU node is a pair of GPUs under one switch in a machine shared with others. A cell is the unit of reservation, not of scheduling: the scheduler may place two 2-GPU jobs in a 4-GPU cell.

The design rests on one assumption, which the paper calls hierarchical uniform composability: all cells of a level are interchangeable for a tenant that asks for that level, and every cell of a level splits into the same number of cells of the level below. A cluster with several GPU models or interconnects is divided into pools that each satisfy it.

How many cells of which level a tenant gets is decided outside HiveD. It depends on budget, business priority and workload. The only technical condition is that the assignment is feasible: all reserved cells of all tenants must fit in the physical cluster at the same time. In the example they fit exactly: 7 + 7 + 18 GPUs in a rack of 32.

Buddy cell allocation

The cells of a VC are logical. A logical cell is bound to a physical cell when a job starts to use a GPU in it, and unbound when none of its GPUs is in use. The binding is dynamic on purpose. It lets HiveD avoid a physical cell with failing hardware, avoid cells that low-priority jobs are using, and pack bound cells together.

The algorithm that manages the binding keeps, for every level, a list of free physical cells (Algorithm 1). At the start only the top level has free cells.

  • Allocate. To hand out a cell of level k, take a free one from that level's list. If the list is empty, allocate a cell of level k + 1 in the same way, split it into its buddies, put them on the free list of level k and take one.
  • Release. When a cell is released, check its buddies. If all of them are free, merge them into the cell of the level above and release that cell in the same way. Otherwise put the cell on its level's free list.

Free capacity is therefore always recorded at the highest level possible, and as many large cells as possible stay whole. Before every allocation HiveD checks that the request is legal: the tenant must hold fewer cells of that level than it reserved.

The figure runs the algorithm on two nodes that three tenants have reserved in full: four requests, one request beyond a reservation, then four releases.

The guarantee is Theorem 1: under hierarchical uniform composability, and if the initial assignment is feasible, buddy cell allocation satisfies every legal request. The proof maintains one inequality for every level k at all times:

rk - ak \le fk + Fk

On the left are the cells of level k still owed to tenants: rk reserved minus ak already allocated. On the right are the cells that can still be provided: fk on the free list, plus Fk, the cells obtainable by splitting higher cells that are not themselves owed to anyone. Every allocation, split and merge keeps the inequality true. The table in the figure shows the two sides at each step.

Because of the uniformity assumption the algorithm never has to look ahead. It does not check after a split whether later requests can still be met. Checking that each request is legal is enough. An allocation or release walks up the levels at most once, so it costs time proportional to the number of levels, which is about five from GPU to rack.

The name comes from buddy memory allocation. The paper states what is new beside the analogy: modeling GPU affinity as cells so that a buddy scheme applies at all, proving a safety property that memory allocators do not need, and replacing the power-of-two rule of memory buddies with a split factor per level.

Low-priority jobs and dynamic binding

Reserved cells that sit idle would be wasted. HiveD lets jobs run at low priority on GPUs that nobody has claimed, and stops them when the owner of a cell needs it. It keeps two views of the same physical cells, one for guaranteed allocations and one for low-priority ones, and runs the same algorithm in both. Preemption has a cost: a stopped job loses the work since its last checkpoint.

Two placement rules keep preemption rare. A low-priority cell is chosen as far as possible from cells in guaranteed use, for example not as a buddy of one. A guaranteed cell is bound to the free physical cell with the fewest GPUs in low-priority use. The second rule is only possible because the binding is dynamic.

The low-priority capacity is divided among tenants by weighted max-min fairness. The paper notes that other fairness schemes could be used instead.

The system

HiveD is about 7,700 lines of Go, open source, and part of OpenPAI, Microsoft's training platform on Kubernetes (Section 4). It runs as a scheduler extender: a separate process that works next to the default Kubernetes scheduler and reuses its basic logic.

  • Cell specification. The operator writes a YAML file with three parts: the cell hierarchy of each kind of hardware with its split factors, the physical cluster as a list of top-level cells and their addresses, and the cells assigned to each VC. HiveD ships a tool that detects infeasible assignments in the file.
  • Faulty hardware. When several free cells are available, HiveD binds to a healthy one. When a VC has no other choice, it binds to the faulty cell anyway, so that the scheduler inside the VC sees the broken hardware and avoids it.
  • Recovery. The free lists and the allocation list live in memory. The binding decision of every pod is stored in the pod's annotation, which Kubernetes keeps reliably, and after a crash HiveD rebuilds its state from the annotations of the running pods.
  • Reconfiguration. A changed specification is handled like a recovery. Jobs whose bindings no longer fit the new assignment are downgraded to low priority and preempted when necessary.

At the time of writing HiveD had run for more than 12 months on a cloud cluster of 800 GPUs. To test failure handling, the authors ran it on 800 GPUs in 200 low-priority Azure VMs, which the cloud can take away at any time. HiveD treats a preempted VM as a faulty cell. The paper reports a preemption rate of up to 75% of the VMs, and that a waiting job was scheduled within a minute of its VM coming back.

Evaluation

The trace

The evaluation uses a two-month trace from a production cluster of 279 8-GPU nodes, 2,232 GPUs, shared by 11 tenants (Table 1). Each of its 141,950 jobs has a submission time, a training time, a number of GPUs with an affinity requirement, and a tenant.

The mix of job sizes differs widely between tenants, and it decides who is exposed to the anomaly. Tenant prod-a asks for 8 GPUs in 72% of its jobs. Tenant prod-e asks for a single GPU in 94% of its jobs, and a 1-GPU job fits anywhere.

Three schedulers, with and without HiveD

The first experiment runs on a real cluster of 96 GPUs in Azure with a ten-day part of the trace, sampled down in proportion to the cluster size (Section 5.1). Three schedulers are tested: a YARN capacity scheduler that packs jobs closely, Gandiva, which migrates jobs to improve placement, and Tiresias, which favors short and small jobs. Each one runs the jobs three ways: every tenant in a private cluster as large as its quota, all tenants sharing the cluster under quota, and all tenants sharing it under HiveD with node-level cells.

All three schedulers show the anomaly under quota, each for a reason tied to its own goal. The capacity scheduler packs jobs, but 1-GPU jobs of varying length from other tenants still fragment the nodes. Gandiva's migrations can improve one tenant's jobs at the expense of another tenant's affinity. Tiresias favors small jobs, and small jobs fragment. With HiveD under them, prod-a's multi-GPU jobs start at once in the periods where they waited before. Average job completion time over all tenants is at most 3% worse, with the capacity scheduler, and up to 12% better, with Gandiva (Figure 5b).

The full trace in simulation

The remaining experiments simulate the capacity scheduler. The simulator is first checked against the real cluster: queuing delays differ by at most 7% (Section 5.2).

  • At the original size, 279 nodes, prod-a queues 200 minutes longer under quota than in its private cluster around days 11 and 12. In that period no node in the cluster can offer 8 free GPUs to a guaranteed job.
  • Under higher load, the same workload on 200 nodes with about 90% of the GPUs in guaranteed use, the anomaly for prod-a lasts from day 9 to day 19, with queuing delays up to 8,000 minutes longer than in its private cluster. Averaged over the trace, prod-a queues nearly 7ร— longer under quota than in its private cluster.
  • Under HiveD, every tenant queues less than in its private cluster. Against quota, HiveD shortens the queuing delay of 9 of the 11 tenants, by up to 94% for prod-a.

The paper then asks what happens if the victim leaves. With prod-a and its GPUs removed, the largest tenant, res-f, becomes the victim and queues more than 1.7ร— longer than it would alone. With res-f removed as well, no anomaly remains, but the two tenants took 37% of the cluster's GPUs with them, and everyone left queues longer than before. The anomaly gives the largest contributors a reason to leave, and their leaving makes the shared cluster worse for those who stay.

Job shapes decide the damage

To isolate the cause, the simulation rewrites the trace. Each tenant's jobs become either all 8-GPU or all 1-GPU, with the tenant's total GPU demand kept the same. At first only res-f runs 8-GPU jobs, and other tenants are moved to 8-GPU jobs one at a time.

The more of the workload consists of 1-GPU jobs, the worse a tenant with 8-GPU jobs fares under quota. Under HiveD the same tenant is never worse off than in its private cluster.

A soft affinity requirement does not remove the problem. With half of the multi-GPU jobs allowed to run in a relaxed shape, and with the assumption that relaxed jobs train at full speed, prod-a still queues 1.3ร— longer under quota than in its private cluster (Figure 9). Relaxed jobs can also fragment the cluster further, which lengthens the wait of the jobs that cannot relax.

The allocator itself

Three results concern buddy cell allocation directly (Section 5.3):

What was measuredResult
GPUs preempted, dynamic against static binding55% fewer
Fragmentation, multi-level against node-level cells10% to 20% lower for most of the trace
Time per allocation on a simulated 65,536-GPU cluster2.18 ms on average

The fragmentation result is a recommendation in disguise. If two tenants each reserve only node-level cells, their two 1-GPU jobs land on two different nodes. If each reserves a level-1 cell for such jobs, the allocator can make the two cells buddies on one node. The paper estimates that reserving cells that match the job sizes keeps roughly 30 more nodes whole in the 279-node cluster. Of the 2.18 ms, 88% goes into ordering cells by their low-priority use.

Critical analysis

Strengths:

  • The problem is real and newly named. The anomaly comes from user complaints in production, and the definition, never worse than a private cluster of the same size, is one a tenant can check.
  • The guarantee is proved, not tuned. Sharing safety holds for any workload and any scheduler on top, as long as the assignment is feasible and the cells are uniform.
  • It separates two concerns cleanly. Three schedulers with different goals run unchanged in their logic and keep their job completion times.
  • It is deployed and open. The code is public, runs under Kubernetes, and handles failures, reconfiguration and preemptible machines.

Limitations:

  • The real cluster is not a cluster of 8-GPU nodes. The 96 GPUs are 24 cloud VMs with four K80 GPUs each, and two neighboring VMs are treated as one 8-GPU node. The node-level affinity in that experiment is a convention of the setup.
  • The jobs are stand-ins. The authors had no access to the code and data of the traced jobs. They run 11 public models in a 6:3:1 mix of language, speech and vision taken from the Gandiva paper, and fast-forward through the iterations between scheduling events.
  • Two baselines were changed by the authors. The capacity scheduler got a different preemption policy, with the remark that the baseline would otherwise be much worse, and Tiresias got a quota mechanism it did not have.
  • The largest numbers come from simulation. The 8,000 minutes, the 7ร— and the 132ร— are simulated with one scheduler, the last one on a rewritten trace. The measured anomaly on real hardware is at most 1,000 minutes.
  • The hard decision is left out. How many cells of which level each tenant reserves is called a business process. The paper's own experiment shows that this choice moves fragmentation by 10% to 20%, and a tenant that reserves only node-level cells for a workload of 1-GPU jobs holds on to whole nodes it does not need.
  • Utilization is not evaluated. The paper reports that HiveD's utilization is similar to quota's or slightly better, by up to 20% at some instants, and states that the evaluation does not focus on it. Reserved cells that are idle are useful only through low-priority jobs, which can lose work when preempted.
  • The deployment has no measurements. The 800-GPU cluster is described by its uptime and by how preempted VMs are handled, not by queuing delays or utilization.
  • Uniformity is assumed. A cluster with mixed hardware is cut into uniform pools, and the guarantee holds inside a pool. The paper does not discuss jobs that could run in more than one pool.

Where it fits

The idea is older than GPUs. Buddy allocation comes from memory management, and reserving a shape instead of a count is what an HPC user does by hand when asking Slurm for whole nodes, or for GPUs on one socket of a NUMA machine. HiveD's contribution is to make that reservation a property of the tenant, to enforce it across tenants, and to prove that it holds.

It also explains a requirement that the training papers take for granted. Megatron-LM keeps a model-parallel group inside one server because it exchanges data several times per layer, and the scheduler described in the NCCLX paper places consecutive ranks as close together as it can. Both need a cluster that can still offer whole nodes when a job asks for them. Sharing safety is the promise that it can.

If you found this paper review helpful, consider sharing it with others.

Mastodon