The previous post made GPU servers into checkpoint caches and made loading from those caches fast. This post is about a problem that follows from it. Once weights live on specific servers, the scheduler wants to send requests thereand the GPU there is frequently busy running something else for an unknown amount of time.

Two servers, two models, one request

The paper’s Figure 3 is a small example that I think is the best way into the problem, so I will go through it.

Two servers, two models. Server 1 has model A in DRAM, model B on SSD and an idle GPU. Server 2 has model B in DRAM and its GPU is busy running an inference of model A. A request for model B arrives. Where does B start?

Server 1: idle GPU, A in DRAM, B on SSD. Server 2: GPU running A, B in DRAM. A request for B arrives.

The paper's Figure 3 setup. The timelines below are redrawn from it; bar lengths are illustrative.

Availability-driven. Take the free GPU: server 1. B loads from SSD, which is slow. A is unaffected. B’s user has to wait.

Timeline: S1 loads B from SSD then infers; S2 keeps running A then idles

B pays a full SSD load.

Locality-driven. Take the server where B is in DRAM: server 2. B loads fast, but only after A finishes and A’s duration is unknown. B queues and server 1 sits idle. This is what locality-aware model servers such as Clockwork and Shepherd do and it works well for models whose inference time is predictable. LLM inference time depends on the output length, which nobody knows in advance.

Timeline: S1 idle; S2 finishes A, loads B from DRAM, infers B

B queues behind an inference of unknown length while S1 sits idle.

Preemption-driven. Remove A from server 2, start B there, restart A on server 1. B starts quickly. However, A has to be loaded on server 1, its KV cache has to be recomputed from scratch and the user of A sees a long pause. Shepherd does this and in the paper’s experiments it is the policy that hurts most under load.

Timeline: S2 evicts A and starts B; S1 loads A and recomputes its KV cache before resuming

B starts quickly, but A is delayed.

Locality-driven with live migration. Load A on server 1 while A keeps generating on server 2. When server 1 is ready, move A’s inference state over, continue there and start B on server 2 from DRAM. Both models get low latency. Neither user sees a long stall.

Timeline: S1 loads A while S2 keeps generating A; A moves to S1, then S2 loads B from DRAM

Only this policy keeps both latencies short.

The fourth policy is the one we want. The open question is how much it costs to move A’s inference state.

Behind the paper: why we did not use snapshots

The obvious approach to live migration is the one that works for virtual machines: snapshot the process and ship the snapshot. Microsoft’s Singularity does checkpoint-and-restore for GPU jobs this way. We looked at it first.

It was much too slow. Creating the snapshot of a GPU inference process and transferring it takes tens of seconds or minutes, which is the latency we were trying to eliminate. The VM world’s refinement, dirty-page tracking so you only ship what changed, is not supported for GPU-enabled containers and VMs. So we could not build on the infrastructure and had to migrate at the application level.

This constraint turned out to be useful, because the state of an LLM inference can be described very compactly.

Tokens are the state

An application-level migration has two objectives. The state you move must be small, to keep network traffic down. And the destination must catch up with the source quickly, so the migration converges.

Small state. The state of an autoregressive inference is the prompt tokens plus every token generated so far. That is tens to hundreds of kilobytes. The KV cache that the GPU has built from those tokens is 1 to 10 GB or more. Instead of sending the cache, we send the tokens and recompute the cache on the destination. Under some conditions, a fast network and a short sequence, shipping the cache could be competitive, but it always adds cluster traffic and it is never smaller.

Fast catch-up. Recomputing the KV cache for a sequence of tokens is a prefill: one forward pass over all of them, batched, using the GPU well. Generating those same tokens was decode: one forward pass per token, each of them bandwidth-bound. Prefill is roughly an order of magnitude faster per token. Recomputing the cache for 1,000 tokens costs about as much as generating 100 new ones. So if the source keeps generating while the destination recomputes, the destination catches up quickly and the gap between them shrinks each round.

I still think this is the most elegant idea in the paper. A VM has no compact description of its state, but an LLM inference does.

The multi-round protocol

The paper’s Figure 4, in three phases. First, prepare the destination:

Sequence diagram: the loading scheduler tells the destination to load A, then tells the source to migrate

Steps 1 and 2. Nothing has moved yet.

  1. The loading scheduler tells the destination to load model A. If an idle instance of A already exists there, skip.
  2. Once loaded, the scheduler sends the source a migrate request carrying the destination’s address.

Then move the inference, in rounds, while the source keeps generating:

Sequence diagram: the source sends tokens to the destination, which recomputes the KV cache while the source keeps generating; repeated

Steps 3 and 4 repeat until the token gap is small. Only tokens cross the network.

  1. The source marks itself “migrating” and, if the inference has not finished, sends the destination a resume request with the intermediate tokens: the prompt plus everything generated so far.
  2. The destination recomputes its KV cache from those tokens. Meanwhile the source keeps generating.

Finally, hand over:

Sequence diagram: the source stops and returns all tokens to scheduler and router; the scheduler unloads A at the source; the router redirects to the destination

Steps 5 to 7: the router switches to the destination and the client keeps receiving tokens.

  1. When the resume completes, the source stops, returns to the scheduler and replies to the request router with all tokens, including the ones it generated during step 4 and a flag that says “migrated”.
  2. The scheduler finishes the migration, unloads A from the source and starts loading B there.
  3. The router sees the flag, swaps the source for the destination in its route table and sends all tokens to the destination, which continues the inference.

The client sees a stream of tokens that pauses for one short round and continues. It never learns that the GPU changed.

Corner cases we had to handle

Because inference length is unpredictable, the source may finish between steps 3 and 5. Then it tells the router the inference is complete as usual and tells the scheduler, which instructs the destination to stop resuming and abandon the migration.

Source failure during loading, before step 2: the scheduler aborts and unloads the destination. Source failure during migration, between steps 2 and 3: the scheduler tells the destination to clear its partial KV cache and unload. Destination failure during loading: cancel. Destination failure during resume: the source hears about it and simply carries on generating where it is.

None of these cases is unusual, but handling them is what makes the mechanism usable in a real system and it is why the migration path is implemented as a state machine rather than a single function call.

Evaluation

The direct comparison is against Shepherd, modified to use our loading-time estimator so that it picks the same GPU we do. The only difference left is that Shepherd* preempts and we migrate. On a four-node cluster with OPT-6.7B:

  • At low request rates there is no locality contention and the two match.
  • At medium rate with ShareGPT, where inferences run 3.7x longer than GSM8K, Shepherd*’s P99 latency is 2x ours because of preemption, even though we migrate more often than it preempts: 114 migrations versus 40 preemptions out of 513 requests.
  • At high rate with GSM8K we beat Shepherd* by 1.27x and a plain serverless scheduler by 1.95x on P99, with 53 migrations against 9 preemptions out of 925 requests. Shepherd* also reads checkpoints from SSD twice as often as we do, because every preemption is a reload.

Migration is cheap enough to do often, while preemption is expensive enough to affect tail latency even when it happens rarely.

One smaller measurement is also worth mentioning: the GPU-side estimate of resume time is accurate to about 5 ms, but we saw an outlier of 623 ms in one of 119 migrations, caused by the CUDA driver taking its time to clean up GPU state when a model is unloaded. The scheduler’s estimates are accurate in general, but the driver can still cause unexpected delays.

What migration needs from the scheduler

Migration is a mechanism. The policy that decides when to use it needs two numbers per candidate server: how long loading the model here would take and how long migrating whatever is running here would take. Getting those numbers and choosing the server with the smallest startup time, is the next post.