Serverless LLMs · part 2 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
Part 1 ended with a baseline: torch.load moves a 2.05 GiB checkpoint to the GPU 1.4x slower than read(2) moves the same bytes into a buffer and three possible causes. This post works through the first three optimisations. Each one ends with a measurement and in one of them I almost published a wrong number.
Optimisation 1: safetensors, measured carefully
The first cause was deserialisation: executing a pickle and rebuilding 201 tensor objects. safetensors addresses exactly that. The header is JSON, the tensor bytes are flat and there is no code to run. So I swapped the loader:
What mmap actually does
opening the file and fetching 'lm_head.weight' (125.00 MiB): 1.0 ms
...because nothing was read. The tensor is a window onto a mapping.
One millisecond for 125 MiB is 125 GB/s, which this laptop cannot do. load_file returns a memory mapping. The bytes get read later, when something touches them, one page fault at a time. Timing load_file therefore measures a syscall, not a load. To measure the real cost, I made the loader touch every byte:
Optimisation 1: load_file() -- and the same load, honestly
stage time share
--------------------------------------------------
load_file (mmap) 2.1 ms 0.1%
touch every byte 1.778 s 62.8%
copy host -> device 1.050 s 37.1%
--------------------------------------------------
TOTAL 2.830 s 100.0%
load_file() reports 2.1 ms (1060.29 GB/s -- physically impossible)
bytes really cost 1.780 s ( 1.24 GB/s)
That is the pessimistic path, faulting through a mapping page by page. The optimistic path is load_file(device="mps"), which lets safetensors copy straight to the device:
safetensors load_file(device='mps') x5
run 1: 686.7 ms ( 3.20 GB/s)
run 2: 540.9 ms ( 4.07 GB/s)
run 3: 561.7 ms ( 3.92 GB/s)
run 4: 576.9 ms ( 3.81 GB/s)
run 5: 569.1 ms ( 3.87 GB/s)
median 569.1 ms min 540.9 ms max 686.7 ms sd 51.2 ms max/min 1.27x
This is faster than the pickle, but the spread is large: 1.27x between the best and worst run of an identical workload. That variance is the page cache. When the file is resident, we measure DRAM. When it is not, we measure the SSD through the fault handler, 16 KiB at a time, with a read-ahead policy the kernel chose. A serverless platform cannot guarantee either state. Users experience the p99 latency and the p99 comes from this spread.
With mmap the kernel decides when bytes move, 16 KiB per fault.
Measuring the copy on the laptop
The second cause was the copy: 201 separate host-to-device transfers. Before trying to fix it, I measured it:
Probe: is the copy bandwidth-bound or call-bound?
201 separate copies : 132.7 ms ( 16.58 GB/s)
1 contiguous copy : 204.4 ms ( 10.76 GB/s)
batching wins : 0.65x (-357 us of overhead per tensor)
Batching the copies does not help on this machine and it is worth understanding why. An M1 has unified memory. There is no PCIe bus between “host” and “device”; the copy is a memcpy inside one DRAM pool. Per-call overhead is invisible next to that. On an A100 the same probe shows the opposite, because every copy there crosses a real bus and the paper measures a 1.4x gain from pinned-memory DMA alone. The lesson for me was to trust measurements on the machine I am using rather than results reported for other hardware. On this Mac the copy stage is not the problem. The remaining time goes into reading bytes from the SSD and processing them on the way.
(Best of three runs each, with the other variant’s buffers freed first. Without that step the benchmark measures swap rather than copies, a mistake I made twice.)
Summary
safetensors is faster than the pickle for good reasons: no code execution, no 201-way unpickling and direct offsets. However, its main benefits are safety and lower variance rather than speed. For this reason the next step aims at predictable load times rather than faster ones, by taking the page cache out of the read path.
Optimisation 2: take the page cache out of the loop
On Linux this is O_DIRECT: the read bypasses the page cache and lands in your buffer, with the constraint that buffers and offsets are 512 B or 4 KiB aligned. macOS has no O_DIRECT. It has fcntl(fd, F_NOCACHE, 1), which has no alignment requirement and tells the kernel not to retain the pages. Enabling it takes one call:
def open_nocache(path: str) -> int:
fd = os.open(path, os.O_RDONLY)
fcntl.fcntl(fd, F_NOCACHE, 1) # F_NOCACHE = 48
return fd
def read_chunks_into(fd, mv: memoryview, chunk: int, offset: int = 0):
done = 0
while done < len(mv):
n = min(chunk, len(mv) - done)
done += os.preadv(fd, [mv[done:done + n]], offset + done)
return done
Now the loader chooses the chunk size instead of the fault handler. I swept it:
Chunk-size sweep over 2.05 GiB (3 runs per point)
page cache
chunk median bandwidth min max max/min
64.00 KiB 340.9 ms 6.45 GB/s 328m 366m 1.12x
256.00 KiB 298.7 ms 7.37 GB/s 288m 319m 1.11x
1.00 MiB 300.7 ms 7.32 GB/s 282m 308m 1.10x
16.00 MiB 300.0 ms 7.33 GB/s 286m 309m 1.08x
64.00 MiB 305.6 ms 7.20 GB/s 298m 306m 1.02x
F_NOCACHE
chunk median bandwidth min max max/min
64.00 KiB 358.6 ms 6.14 GB/s 355m 361m 1.02x
256.00 KiB 327.6 ms 6.72 GB/s 316m 337m 1.07x
1.00 MiB 245.4 ms 8.97 GB/s 223m 310m 1.39x
16.00 MiB 281.7 ms 7.81 GB/s 271m 283m 1.04x
64.00 MiB 258.0 ms 8.53 GB/s 221m 261m 1.18x
A number I nearly published
F_NOCACHE promises the kernel will not retain these pages. It does not promise to ignore pages that are already resident and after the first sweep most of this file is. So the second table is an upper bound on what the SSD can do, not a measurement of it. To get a genuinely cold read on macOS you have to evict first with /usr/sbin/purge (root) before every read, which is what the --cold flag in the harness does. I nearly published the 8.97 GB/s number as “the SSD number” before noticing this.
What remains valid either way is the shape of the curve. At 64 KiB every chunk costs a syscall and a device round trip, so the small-chunk rows are slow. From about 1 MiB up the queue stays full. Chunk size matters and it matters in the same direction hot or cold.
The two tables answer different questions. The first is the number usually quoted, but it leaves out an important detail: the cached read is fast because the bytes never left DRAM and the result no longer holds once another tenant, another model or a large compile evicts the pages. The second table shows what the hardware delivers every time, regardless of what else runs on the machine. This is why ServerlessLLM reads checkpoints in large aligned chunks with the cache bypassed instead of using mmap. The reason is not that it is always faster, but that its performance is consistent.
With F_NOCACHE the loader decides when bytes move and in what chunk size.
Optimisation 3: overlapping the read and the copy
If the loader reads everything and then copies everything, each device is idle while the other works:
Serial baseline: read all, then copy all
SSD -> DRAM 239.7 ms ( 9.18 GB/s)
DRAM -> GPU 46.1 ms ( 47.77 GB/s)
total 285.8 ms ( 7.70 GB/s)
For 46.1 ms of that, one of the two devices is idle waiting for the other.
The solution is a pipeline. Split the file into slabs. Give a pool of reader threads a ring of staging buffers. Have the main thread drain the ring into one device buffer while the readers are still going. The ring is bounded, so readers block when the copier falls behind and peak host memory is five buffers whatever the model size.
def reader() -> None:
fd = open_nocache(path) # one descriptor per thread
while True:
try: s = todo_q.get_nowait()
except queue.Empty: return
off, n = s * slab, min(slab, size - s * slab)
b = free_q.get() # wait for a staging buffer
read_chunks_into(fd, views[b][:n], chunk, offset=off)
ready_q.put((b, off, n))
# main thread, concurrently:
for _ in range(nslabs):
b, off, n = ready_q.get()
dev[off:off + n].copy_(staging[b][:n])
free_q.put(b)
Readers fill the ring from one end, the copier drains it from the other. The free queue is the back-pressure.
The “read only” column below is the same reader pool with the copier removed. It is there so the two columns separate “threads made the read faster” from “the copy is now hidden behind the read”:
Pipelined: 64.00 MiB slabs, 4 staging buffers, 16.00 MiB preads
readers read only pipelined bandwidth vs serial copy hidden
1 201.7 ms 281.9 ms 7.81 GB/s 1.01x 0%
2 126.6 ms 198.6 ms 11.08 GB/s 1.44x 0%
4 97.5 ms 136.3 ms 16.14 GB/s 2.10x 16%
8 58.8 ms 147.0 ms 14.96 GB/s 1.94x 0%
Slab-size sweep at 4 readers
slab pipelined bandwidth
16.00 MiB 200.6 ms 10.97 GB/s
64.00 MiB 139.8 ms 15.73 GB/s
256.00 MiB 279.8 ms 7.86 GB/s
The 2.10x combines two separate improvements and I want to separate them before claiming either.
Concurrency on the read. One thread moved the file at 9.18 GB/s. Four threads moved the same bytes at over 20 GB/s, with the GPU not involved at all. One thread cannot keep a storage queue full: it is blocked on the device for most of every pread. The paper’s breakdown finds the same thing on server NVMe, where multiple threads are worth 2.3x on their own because the SSD has multiple channels to feed.
Overlap with the copy. Serial cost 285.8 ms. The read alone at four threads costs 97.5 ms and the pipeline lands at 136.3 ms. So the copy is not fully hidden and the pipeline is limited by the copier. Every slab is a separate device copy issued from one Python thread that holds the GIL between preads and 46 ms of contiguous copy becomes more expensive once it is split into slabs and interleaved with reads. This is a real limitation of doing this in Python and it is the main reason the ServerlessLLM store is written in C++ with its own thread pool per storage tier.
Most of the 2.10× is read concurrency. The gap between the last two bars is copy that did not get hidden.
The reader column stops improving at some point. Threads help until the device is saturated and after that they make things worse: they cannot add bandwidth, only more queueing, which increases the p99 latency. ServerlessLLM applies the same reasoning per tier: it uses enough concurrency to keep each device busy and no more.
The same caveat applies here: with the page cache warm these read numbers are DRAM numbers. 20 GB/s is not what an M1 Pro SSD does. Re-run after a purge to see the cold shape; the ordering of the variants holds.
Where this leaves us
The loader now moves bytes about as fast as this process can. This raises a question: why move 2 GiB through a ring buffer and then rebuild 201 tensors on top of it, when the layout we want is fixed and known before the process even starts?
Part 3 answers this by changing the file format instead of the reader.