Serverless LLMs · part 4 of 7
- Anatomy of an LLM Cold Start, Part 1: Where the Time Goes
- Anatomy of an LLM Cold Start, Part 2: Make It Predictable, Then Make It Fast
- Anatomy of an LLM Cold Start, Part 3: A Checkpoint Format Written for the Reader
- ServerlessLLM Paper Breakdown, Part 1: Your GPU Server Is Also a Storage Server
- ServerlessLLM Paper Breakdown, Part 2: Migrate Tokens, Not Gigabytes
- ServerlessLLM Paper Breakdown, Part 3: Scheduling for Startup Time
- Serverless Agents: When Every Step Is a Cold Start
Why I am writing about my own paper
I like paper-journal posts where someone explains a systems paper through the life of one request. The best ones are written by outsiders who found the paper interesting. In this post I do something different and cover a paper I helped write: ServerlessLLM: Low-Latency Serverless Inference for Large Language Models, OSDI 2024.
The advantage of being a co-author is that I can explain what we tried first and dropped and which parts of the paper matter most. The disadvantage is bias, so I keep the numbers next to the claims and credit the people who did the work. Yao Fu led the project in Luo Mai’s group at Edinburgh and the system is open source.
The first three posts of this series built a fast checkpoint loader on a laptop and reached a limit: loading got 1.6x faster and the user’s wait only fell 1.13x, because a cold start has fixed costs that no file format can reduce. This post is about what a cluster can do about that. It covers the motivation and the first of three designs. The next two posts cover live migration and the scheduler.
Why serverless for LLMs at all
Serverless inference means you upload a model and you pay only while a request is running. Amazon, Azure, Google, Hugging Face, Together, Replicate, Databricks, Fireworks and Cohere all sell some version of it. The appeal is obvious for anything with bursty or unpredictable load: a product that just launched, an internal tool in healthcare or legal that is idle most of the day, a fine-tuned model with a dozen users. Reserving GPUs for these means paying for GPUs that are idle most of the time.
The provider also benefits. If many models can share a GPU pool, utilisation goes up and the provider earns a software premium for running the infrastructure. This works for both sides as long as one condition holds: a model that is not running has to start fast when a request arrives.
In practice this condition often fails. The Azure Functions trace shows that more than 40% of functions have a cold-start rate above 25%and about a quarter of functions above 60%, within a 5-minute keep-alive window. Those are ordinary functions. LLMs are worse, because the thing being started is tens to hundreds of gigabytes. Providers have said in public that initialising a state-of-the-art LLM on their serverless platform takes tens of seconds.
Life of a serverless inference request
To see where those seconds go, I follow one request through a typical GPU serverless cluster.
- A request arrives at the controller. Its request router checks whether an inference process for that model is already running somewhere. If so, the request is routed there and there is no cold start.
- If not, the model loading scheduler picks a node with a free GPU and tells it to start the model.
- The node starts a GPU process or container and sets up an inference library.
- That process downloads the checkpoint from remote model storage onto local disk.
- It loads the checkpoint: initialises the model object, allocates GPU memory, creates tensors, copies bytes. SSD, then DRAM, then GPU.
- Prefill runs the prompt through the model, producing the first token and the KV cache.
- Decode generates one token per iteration until end-of-sequence. Each token is streamed back as it is produced.
Steps 6 and 7 are the inference. Everything before is the cold start. Users measure two things: time to first token, which the cold start dominates and time per token, which depends on the model and the hardware. An inference can run for seconds or minutes and you cannot know which in advance, because the output length is not known until the model emits end-of-sequence. That unpredictability matters a lot in the next post.
Cold starts at scale
Two facts from the paper’s motivation section.
Checkpoints are large, so downloads are slow. Grok-1 is over 600 GB in fp16. DBRX is 250 GB, Mixtral-8x22B about 280 GB. A 130 GB checkpoint like LLaMA-2-70B takes at least 26 seconds to pull from S3 or blob storage over a fast 5 GB/s commodity network and most clusters do not have that.
Loading is slow even from local NVMe. With PyTorch, loading OPT-30B into four GPUs takes 34 seconds and LLaMA-2-70B into eight GPUs takes 84 seconds. This is the part I took apart on a laptop in the earlier posts: initialising the model, allocating GPU memory, creating thousands of tensors, copying them one by one. Generating a token usually takes under 100 ms. So the wait before the first token is as long as generating several hundred tokens.
What people do about it today
Three approaches were common when we wrote the paper and each one has a specific problem at LLM scale.
Over-subscribe GPUs. Keep a pool of warm instances so nobody hits a cold start. This is what AWS Serverless Inference does and it works for ResNet and BERT. For a model that needs four or eight expensive GPUs to run at all, keeping warm spares is exactly the cost that serverless was supposed to avoid.
Cache checkpoints in host memory. Systems like Clockwork keep model weights in DRAM so a start never touches disk. This works for models of a few GB. LLMs are hundreds of GB and the number of distinct models on a serverless platform is large, so host memory alone has frequent misses and each miss means a download.
Add storage servers. Put a fast cache tier inside the cluster. Trace studies show downloads still exceed 20 seconds through an optimised pipeline over a 100 Gbps NIC and the hardware is not cheap. A network-optimised ElastiCache node with 210 GB of memory and 200 Gbps costs about $16 per hour, roughly the price of an 8-GPU g5.48xlarge. Caching a 70B model this way doubles the cost.
The observation
The key observation of the paper is that a GPU inference server already contains a large, fast, multi-tier storage hierarchy and serverless systems barely use it.
Capacity: a contemporary 8-GPU server supports up to 4 TB of DRAM, 64 TB of NVMe SSD and 192 TB of SATA SSD. Bandwidth: each GPU has its own PCIe link to host memory, so an 8-GPU PCIe 5.0 server has about 512 GB/s aggregate between DRAM and GPUs and RAID-0 NVMe gives around 60 GB/s from SSD to DRAM. Compare that with the 10 Gbps, or about 1.25 GB/s, that a typical cluster network provides for downloads.
In a serverless inference context we observed that most of that host memory and most of those SSDs sit idle. Existing systems like KServe and Ray Serve use a fraction of the host memory and minimally use SSDs for caching.
The storage inside a GPU server can be used as a checkpoint cache.
Treating that storage as a checkpoint cache is cost-effective, because it is already there. It is scalable, because every new server adds its own DRAM and SSDs. It also improves with newer hardware: a Grace Hopper node has 1 TB of on-chip DRAM and a 900 GB/s link to HBM.
However, using this storage well is not trivial. Three problems have to be solved and they map onto the three designs in the paper:
- Loading from a multi-tier hierarchy has to be fast and predictable. Checkpoint formats and loaders were built for training, where files are written often and read rarely. In serverless the opposite holds. (This post.)
- Locality only helps if requests actually go where the weights are and the GPU there may be busy with an inference of unknown duration. (Next post: live migration.)
- Someone has to decide, per request, whether loading from SSD here beats moving an inference there and that needs startup-time estimates. (The post after: scheduling.)
Design 1: loading-optimised checkpoints and a multi-tier pipeline
The first three posts arrived at this design on a laptop. The paper’s version is the same idea with real hardware and multiple GPUs.
The format. We assume a checkpoint has execution files, which define the architecture and the model-parallel plan and parameter files, which hold the bytes. We convert the parameter files into one sequential partition per GPU, containing nothing but tensor bytes, plus a tensor index that maps each name to a tuple of GPU id, offset and size. Tensors are aligned to memory word sizes so an address is a base plus an offset. Metadata is kept out of the partitions so reads can be large and sequential.
The multi-tier loading subsystem. A pinned-memory chunk pool sits in DRAM, managed as fixed-size chunks to avoid fragmentation and with explicit allocate and free APIs so the application, not an LRU heuristic, decides what stays. Reads from SSD use direct I/O, bypassing the page cache for predictable performance. Each storage tier, remote object store, local SSD and pinned memory, has its own pool of I/O threads. Threads in one tier read chunks and enqueue chunk indices for the threads in the next tier, so the tiers form a pipeline and all of them stay busy. Copies to the GPU go from pinned memory by DMA, using the parallel PCIe links when a model is partitioned across GPUs.
Each tier has its own thread pool, so the load time is determined by the slowest tier.
Each optimisation contributes. On RAID-0 NVMe with OPT models, going from reading tensor by tensor to the full pipeline breaks down as:
| optimisation | throughput gain |
|---|---|
| bulk reads instead of per-tensor reads | 1.2x |
| direct I/O, bypassing the page cache | 2.1x |
| multiple I/O threads per tier | 2.3x |
| pinned memory, DMA to GPU | 1.4x |
| pipelining across tiers | 1.5x |
Four CPU cores are enough to hit full bandwidth, with 16 MB chunks. The loader saturates what FIO can pull from the same device, from SATA up to RAID-0 NVMe and it keeps scaling on faster devices where PyTorch and safetensors flatten out. End to end it loads checkpoints 3.6x faster than safetensors and 6x faster than PyTorch on OPT-2.7B and 4.7x and 8.2x on LLaMA-2-70B. safetensors takes 112,000 page faults to load LLaMA-2-7B cold; the direct-I/O path takes none.
Behind the paper: the loader is a separate process
I want to highlight one design choice, because it looks like an implementation detail in the paper but has important consequences.
The checkpoint loader is a model manager that runs as its own long-lived process, not a library inside the inference process. The manager allocates GPU memory and streams the bytes in. The inference process, vLLM or Transformers or whatever you run, initialises the model object on its own, then asks the manager for the base address of each GPU allocation as a CUDA IPC handle, reads the offset of each tensor from the index and sets each tensor’s data pointer to base plus offset. A single synchronisation at the end makes sure the bytes have landed before inference starts.
Three things follow from that split:
- Loading can be scheduled before the inference process is ready and overlapped with its initialisation, which was the 2 seconds of Python imports that dominated the laptop measurement in part 3.
- The manager outlives any single worker, so the host-memory copy of a model is shared. A second replica of a model that is already resident never touches the disk.
- The inference framework needs almost no changes. It builds the model on a meta device and gets pointers handed to it.
What fast loading does not solve
Fast loading on one server is necessary but not sufficient. If the scheduler sends a request to a server that has the checkpoint only on SSD, or not at all, the fast loader does not help much. And if it sends the request to the server that has the model in DRAM, that server’s GPU may be halfway through generating a 2,000-token answer for someone else.
The next post is about this problem and about the idea in the paper I like most: when an inference has to move, move its tokens instead of gigabytes of state.