System Design: Training & ML Infrastructure
ML-infra design rounds ask you to build the machinery behind models: the data pipeline that turns a web crawl into training tokens, the scheduler that shares ten thousand GPUs among a hundred researchers, the control plane that keeps a 16k-GPU run alive through a failure every few hours, the RL rollout fleet, the serving platform, the eval farm. The interviewer wants classical distributed-systems engineering (schedulers, storage tiers, queues, metadata stores, idempotency, autoscaling) applied under ML constraints: GPU-hours cost real money, failures are routine, data is measured in petabytes, and every result has to be traceable to the exact data and code that produced it. This page is the design-interview counterpart to A3, A5 and A6. Those pages explain the ML concepts. This one shows how to build the systems around them.
TL;DR: the 8–12 things to be able to say out loud
- The metric is goodput, not uptime. Useful training progress per GPU-hour paid for. Goodput = availability × (1 − checkpoint and restart overhead) × MFU. Every design choice should name which factor it improves.
- At scale, failure is a rate, not an event. Job MTBF ≈ per-GPU MTBF ÷ GPU count. Meta saw about one interruption every 3 hours on a 16k-GPU Llama 3 run. Design for automatic detection, eviction, restart and resume.
- Checkpoint interval is derived, not chosen: \(\tau^* \approx \sqrt{2CM}\) (Young/Daly). Async and in-memory checkpoints shrink \(C\). Hot spares and fast restart shrink the lost time per failure.
- Name the bottleneck before proposing a fix: compute, HBM capacity or bandwidth, interconnect, storage I/O, data loader, or the scheduler. Each has its own symptom and its own metric.
- Schedulers for ML need gang scheduling, topology awareness, quotas with borrowing, and preemption. Slurm has had these for years. On Kubernetes you add them with Kueue or Volcano.
- Storage is tiered: object store (durable, cheap, high aggregate throughput) → parallel FS or cache (shared POSIX, fast) → local NVMe (per node, fastest). Training reads sequential shards, never millions of small files.
- Web-scale data pipelines are big shuffles. MinHash-LSH dedup turns into a group-by on band keys plus connected components. Plan for reprocessing whenever filters change, and record lineage for every shard.
- RL infrastructure is an inference fleet attached to a trainer, plus reward services and thousands of sandboxes. The hard parts are weight sync, staleness, long-tail rollouts, and keeping sandboxes warm.
- Serving platforms are dominated by cold start and routing. Weights take tens of seconds to minutes to load. Autoscale on queue depth or KV-cache utilization, not CPU. Route by prefix-cache affinity.
- Always do the math out loud: GPU-hours, bytes moved, time per stage, dollars. A design with numbers attached is easier to defend than a design with boxes only.
Part 0 · How ML-infra design interviews differ
A classic system design round ("design Twitter") is about request paths, read/write ratios, caching and consistency. An ML-infra round keeps all of that and adds a different economics. A single H100-class GPU rents for roughly $2–4 per hour. A 16k-GPU cluster therefore burns about $30–65k per hour whether or not it is doing useful work. The DeepSeek-V3 report priced its final training run at 2.788M H800 GPU-hours, assuming $2 per GPU-hour DeepSeek-AI 2024. When idle hardware costs that much, the important questions are about efficiency and recovery. How much of each paid GPU-hour turns into useful progress? How fast do you recover when a component fails? Can you prove which data produced which model?
The core metrics
| Metric | Definition | What moves it |
|---|---|---|
| MFU (model FLOPs utilization) | Achieved model FLOPs ÷ peak hardware FLOPs while the job is running. See A3. | Parallelism strategy, kernels, communication overlap, batch size |
| Goodput | Useful training progress (tokens or steps that survive into the final model) ÷ what ideal, uninterrupted hardware would deliver in the same wall time | Failure rate, detection time, restart time, checkpoint overhead, lost work since the last checkpoint, stragglers |
| Effective training time | Fraction of wall-clock time spent training, not restarting or recovering. Llama 3 reported above 90% Llama Team 2024 | Same as goodput, minus the MFU term |
| Allocation vs utilization | Allocation: GPUs assigned to jobs ÷ GPUs owned. Utilization: SM activity or MFU of those allocated GPUs | Scheduler policy, fragmentation, idle interactive sessions, data-loader stalls |
| GPU-hours and $ | The budget unit for every job | All of the above, plus hardware choice, spot or preemptible capacity, right-sizing |
Treat goodput as a product of factors: \(\text{goodput} \approx A \times (1 - o_\text{ckpt}) \times (1 - o_\text{lost}) \times \text{MFU}\), where \(A\) is the fraction of time the job holds healthy hardware, \(o_\text{ckpt}\) is the time training is blocked by checkpointing, and \(o_\text{lost}\) is the recomputed work after a rollback. A 5-point gain in any factor is worth the same. That is why infra teams care about restart time as much as kernel teams care about MFU.
Failures are the normal case
With \(N\) GPUs that fail independently, job MTBF ≈ per-GPU MTBF ÷ \(N\). In a 54-day snapshot of Llama 3 405B pre-training on 16k H100s, Meta saw 466 job interruptions, 419 of them unexpected. About 78% of the unexpected ones were attributed to confirmed or suspected hardware issues, with faulty GPUs and HBM3 memory the largest categories Llama Team 2024. 419 ÷ (54 × 24 h) is about one unexpected interruption every 3.1 hours. Working backwards, that implies a per-GPU interruption rate of roughly once every ~50k GPU-hours (≈ 6 years). Each GPU is quite reliable. Sixteen thousand of them, all in one synchronous job, are not.
Data volumes and the shape of I/O
Pre-training corpora are tens of trillions of tokens. FineWeb alone is 15T tokens from 96 Common Crawl snapshots Penedo+ 2024, and DCLM's raw pool is 240T tokens Li+ 2024. At ~4 bytes per token that is tens of TB of clean text and PBs of raw crawl. Checkpoints are TBs each. Text pre-training has a counterintuitive property, though. The loader's bandwidth need is tiny. 16k GPUs consuming ~2M tokens/s read only about 8 MB/s of tokenized data. The hard parts of a text loader are determinism, resumability, and mixing, not bandwidth. Checkpoints, video and RL environments are where storage bandwidth starts to matter.
The bottleneck taxonomy
| Bottleneck | Symptom | How you see it | Typical fixes |
|---|---|---|---|
| Compute | High SM activity, MFU near the hardware's realistic ceiling | Profiler shows GEMM-dominated timeline | Lower precision (FP8), better kernels, more GPUs |
| HBM capacity | OOM, forced small micro-batches, heavy activation recompute | Memory snapshots, allocator stats | Sharding (FSDP/ZeRO), activation checkpointing, offload, more TP/PP (A3) |
| HBM bandwidth | Low arithmetic intensity, e.g. decode | Roofline position, DRAM throughput counters | Bigger batches, quantization, KV-cache compression (A6) |
| Interconnect | GPUs idle inside collectives. Step time grows with scale | NCCL timings, exposed communication in traces, link counters | Overlap comm with compute, topology-aware placement, change the parallelism layout |
| Storage I/O | Slow checkpoint save or load, slow startup | FS or object-store throughput, time-to-first-step | Async or sharded checkpoints, local caches, parallel reads |
| Data loader | GPU gaps between steps, CPU pegged on preprocessing | Data-wait time per step, host CPU usage | Pre-tokenize, more workers, prefetch, move preprocessing offline |
| Scheduler and queue | Jobs wait while GPUs sit idle, or big jobs never start | Queue time, allocation rate, fragmentation | Gang plus backfill, topology-aware packing, preemption, quotas with borrowing |
| Stragglers | Step time set by the slowest rank. p99 rank time ≫ p50 | Per-rank step timing, NCCL wait attribution | Detect and evict slow nodes, fix thermals or links |
An answering framework
Strong candidates split the control plane from the data plane early. Every ML system here has a small, strongly consistent control plane (job specs, leases, state machines, metadata, in Postgres, etcd or a similar store) that orchestrates a large, failure-prone data plane (GPU workers, storage, network). Most failure-handling questions get easier once you can say "the controller holds the desired state and reconciles toward it, and the workers are cattle."
Numbers to know
All values are approximate and vary by SKU, generation and configuration. Check vendor datasheets before quoting them. Dense means without structured sparsity, since datasheets often headline the 2:4-sparse figure.
| Quantity | Ballpark | Use it for |
|---|---|---|
| H100 SXM BF16 dense | ~990 TFLOPS (FP8 ~1,980) | Training time from 6ND. Realistic MFU 35–50% |
| H100 HBM | 80 GB, ~3.35 TB/s | Decode speed, KV capacity |
| B200-class GPU | ~2.2 PFLOPS BF16 dense, ~180 GB HBM3e, ~8 TB/s | Next-gen capacity planning |
| NVLink (per GPU, total) | H100 ~900 GB/s. Blackwell ~1.8 TB/s | TP inside a node or NVL72 domain |
| InfiniBand / RoCE per NIC | NDR 400 Gb/s ≈ 50 GB/s. 8 NICs per node ≈ 400 GB/s/node | DP, PP, weight broadcast across nodes |
| PCIe Gen5 x16 | ~64 GB/s | Host↔GPU copies (async checkpoint staging) |
| Local NVMe | ~5–12 GB/s per drive | Shard caches, staged checkpoints, model weights |
| Parallel FS / storage cluster aggregate | Hundreds of GB/s to several TB/s. Llama 3's storage: ~2 TB/s sustained, ~7 TB/s peak | Checkpoint save/restore time |
| Object store | Per-prefix request limits (S3: 3,500 PUT / 5,500 GET per second per prefix). Throughput scales with parallel connections. ~100 MB/s per stream | Datasets, durable checkpoints, model registry |
| Training state bytes per param | BF16 weights 2 B. Mixed-precision Adam training state ~14–16 B | Checkpoint size: 70B ≈ 1 TB, 405B ≈ 6 TB, 1T ≈ 15 TB |
| Training FLOPs | ≈ 6 × params × tokens (dense) | GPU-hours = FLOPs ÷ (peak × MFU × 3600) |
| Inference FLOPs | ≈ 2 × active params per token | Batch-job and prefill cost |
| Job MTBF at scale | Llama 3: ~3 h at 16k GPUs. Scales ~1/N | Checkpoint interval, spare pool size |
| Bytes per token (text) | ~4 bytes raw text. 4 bytes stored as uint32 token IDs | Corpus sizing |
| GPU price | ~$2–4 per H100-hour. Spot or preemptible often 30–70% cheaper | Cost estimates |
| Spot / preemptible notice | AWS spot: 2-minute warning. GCP Spot VMs: best-effort shutdown period of up to 30 s (a longer preemption notice is configurable) | Shard sizing for batch jobs |
Hardware generations turn over every 1–2 years, and cloud GPU prices have moved a lot since 2023. Treat this table as a way to get the order of magnitude right in an interview, not as a spec sheet.
- The Llama 3 Herd of Models, §3.3: infrastructure, scaling and reliability statistics from a 16k-GPU run.
- MegaScale (ByteDance, 2024): production lessons from training on 10k+ GPUs, including fault tolerance and diagnosis tooling.
- The Ultra-Scale Playbook: practical parallelism and throughput math.
- A3 · Distributed Training: MFU, memory accounting, parallelism. This page assumes it.
Part 1 · Building blocks
1.1 Cluster schedulers
The scheduler decides which job gets which GPUs and when. ML workloads break the assumptions of general-purpose schedulers in four ways:
- Gang scheduling. A 512-GPU job needs all 512 at once. Starting 500 of them wastes 500 GPUs waiting for the rest, and two half-placed jobs can deadlock each other. All-or-nothing admission is required.
- Topology awareness. Bandwidth inside an NVLink domain is far higher than across the network, and bandwidth under one leaf switch is higher than across the spine. Placing a TP group across nodes, or a DP group across spine layers, can cost double-digit percentages of throughput.
- Preemption with checkpoint semantics. Evicting a training job is fine if it checkpointed recently. Killing it 50 minutes after its last checkpoint wastes 50 minutes × N GPUs.
- Quotas, priorities and fair share across teams, with borrowing, so idle quota doesn't sit unused.
Slurm
The HPC incumbent. It has a batch-job model, partitions and QOS, a multifactor priority plugin with fair-share (Fair Tree), backfill scheduling, preemption by partition priority or QOS Slurm docs, and topology-aware placement via tree or block topology plugins Slurm docs. Multi-node jobs are allocated all-or-nothing by design.
Kubernetes + Kueue / Volcano
The default Kubernetes scheduler has historically placed pods one at a time, with no notion of a gang or a queue. Native gang-scheduling support has been under development upstream, so check its current status. Kueue adds job-level queueing and quota: ClusterQueues with resource quotas, LocalQueues per namespace, cohorts that let queues borrow each other's unused quota, preemption, fair sharing, and topology-aware scheduling Kueue. Volcano is a batch scheduler with PodGroups for gang scheduling (minMember) and queue policies Volcano. JobSet models multi-role distributed jobs JobSet.
Many large labs use both: Slurm (often on bare metal) for the big pre-training runs, and Kubernetes for everything service-shaped (inference, evals, data pipelines, RL environments). Some run Slurm on top of Kubernetes. In an interview, say that either can work, then justify the choice by workload mix.
Key scheduler concepts
| Concept | What it means | Design note |
|---|---|---|
| Gang scheduling (all-or-nothing admission) | Admit a job only when all its pods or nodes can start together | Without it you get partial allocation and deadlock. Kueue admits whole workloads. Volcano uses PodGroup minMember. Note that Slurm's own "gang scheduling" feature means time-slicing jobs on shared nodes Slurm docs, which is a different idea |
| Topology-aware placement | Pack a job into the fewest switches, NVLink domains or blocks | Encode topology as node labels (rack, leaf, block). Prefer tight packing for large jobs and leave whole blocks free |
Treating GPUs as interchangeable integers. A request for "64 GPUs" really means "8 full nodes, each with all 8 GPUs and all NICs, ideally under the same leaf switch, with healthy NVLink". Scheduling at GPU granularity across nodes produces jobs whose collectives cross the spine and run much slower. Say "node-granular allocation for multi-node jobs, GPU-granular only for small single-node jobs".
1.2 Storage tiers and data loading
- Object storage is the durable layer. It is cheap, scales almost without limit, and has high aggregate throughput when you use many parallel connections. Per-request latency is tens of milliseconds and there are per-prefix request-rate limits: S3 documents 3,500 PUT/COPY/POST/DELETE and 5,500 GET/HEAD requests per second per prefix AWS docs.
- Parallel file systems such as Lustre and WEKA, or DeepSeek's open-sourced 3FS, give POSIX access at very high aggregate bandwidth. They suit checkpoints and shared environments, but per GB they cost several times more than object storage, and metadata servers become a bottleneck under small-file storms (for example, 10k ranks all calling
stat()at startup). - Local NVMe is the cache. Stage dataset shards and model weights there. Write checkpoints there first, then upload.
Streaming dataset formats
The standard pattern for training data is to pre-tokenize into large shards (hundreds of MB each), store them in object storage, and stream them with a deterministic, resumable sampler. Two widely used open-source formats:
- WebDataset stores samples as files inside POSIX tar shards, read sequentially, so each shard is one big streaming read. It works with any storage that can serve a byte stream WebDataset.
- MosaicML Streaming (MDS format) streams shards from object storage to a local cache. It is designed for deterministic sample order regardless of the number of nodes, mid-epoch resumption, and elastic restarts on a different number of GPUs MosaicML Streaming.
For LLM text you will often see something simpler: flat binary arrays of token IDs plus an index of document boundaries (Megatron-style .bin/.idx), memory-mapped from local disk. The sampler is a pure function of (seed, global step, data-parallel rank). Resuming from step \(s\) means recomputing the index, with no saved iterator state. That one design choice gives you reproducibility, elastic resume and cheap skip-ahead after a bad batch.
"How does your data loader resume after a crash?" A weak answer saves the iterator object. A strong answer: the global sample order is a deterministic permutation of (seed, epoch), each rank's slice is a function of (step, rank, world size), so resuming at step \(s\), even with a different world size, needs only \(s\) and the config. Mention that changing the data mixture mid-run means versioning the sampler config too.
1.3 Distributed data processing: Spark, Ray Data, Dask
| Apache Spark | Ray Data | Dask | |
|---|---|---|---|
| Model | Bulk-synchronous stages over partitioned DataFrames or RDDs. Shuffle is a first-class operation | Streaming execution over blocks, with heterogeneous CPU and GPU operators in one pipeline Ray docs | Task graphs over pandas/NumPy-like collections in Python |
| Best at | Huge joins, group-bys, dedup shuffles, SQL over PB-scale tables, mature fault tolerance | Pipelines mixing CPU preprocessing with GPU model inference (classifiers, embedding, batch LLM), and streaming into training | Python-native analytics at medium scale, research workflows |
| Weak at | GPU stages (possible but awkward), Python UDF overhead | For the very largest joins and shuffles, Spark is the more widely deployed path. Benchmark Ray on your own data before committing | Very large shuffles, operational maturity at PB scale |
| Typical use | Dedup, joins with URL blocklists, aggregation | Quality-classifier scoring, batch inference, tokenization feeding training | Ad-hoc dataset analysis |
A common split in LLM data pipelines: Spark (or a purpose-built tool like Hugging Face's datatrove, NVIDIA's NeMo Curator, or DeepSeek's smallpond on DuckDB and 3FS) for the shuffle-heavy dedup and joins, and Ray Data for stages that need GPUs, such as model-based quality filters. The underlying cost to reason about is the shuffle. Any operation that groups by a key other than the partition key (dedup, joins, global sort) moves the whole keyed dataset across the network and usually spills to disk. Estimate the bytes shuffled. That number usually dominates the cost.
1.4 Metadata, lineage and experiment tracking
ML systems produce many artifacts (datasets, checkpoints, eval results, served models), each derived from others. The metadata layer answers: what produced this, and what was produced from it?
| Need | Typical tools | Design core |
|---|---|---|
| Experiment tracking | Weights & Biases, MLflow, in-house | Run ID → config, code commit, metrics time series, artifacts. Metrics go to a time-series store. Run metadata goes to a relational DB |
| Model registry | MLflow Model Registry, in-house | Model name → versions → stage (candidate, staging, prod), with links to training run, eval report and approvals |
| Dataset versioning | Apache Iceberg / Delta Lake snapshots, lakeFS, DVC | Immutable snapshots with IDs. A training job pins a snapshot ID, never "latest" |
| Lineage | OpenLineage, in-house graph | A DAG of (input artifacts, code version, config) → output artifact, emitted by every job |
The principle: artifacts are immutable and content-addressed, or at least snapshot-addressed. Names are mutable pointers. "prod-model" points to version 47, which points to checkpoint s3://…/step_120000, produced by run r-8812 from dataset snapshot ds-2026-09-14-v3 and commit abc123. Everything downstream records the IDs it consumed.
1.5 Queues and messaging
Queues show up in ML infra in three places:
- Work queues for batch jobs (SQS, Pub/Sub, Redis streams, or a Postgres table with
SELECT … FOR UPDATE SKIP LOCKED): work items with leases and visibility timeouts, at-least-once delivery, idempotent workers. This is the backbone of embedding backfills, eval runs and synthetic data jobs. - Event streams (Kafka): high-volume append-only logs such as inference request logs, feature events for recommenders, rollout trajectories, annotation events. Use them when there are many independent consumers or you need replay.
- Request queues inside serving: per-model queues in front of GPU replicas. Their depth is the main autoscaling signal.
Large training jobs themselves do not use message queues on the hot path. Gradients move through NCCL collectives over NVLink and InfiniBand or RoCE. Queues belong to the control plane and the asynchronous pipelines around training.
1.6 Observability for GPU fleets
- Hardware telemetry: NVIDIA DCGM exposes GPU health and performance counters (SM activity, tensor-core activity, memory bandwidth, temperatures, power, ECC and XID errors, NVLink counters). Its active diagnostics can run burn-in tests. dcgm-exporter exports these to Prometheus.
- Communication:
NCCL_DEBUG=INFOand related environment variables NCCL docs show topology and transport choices. PyTorch's Flight Recorder keeps a ring buffer of recent collectives per rank PyTorch docs, so after a hang you can see which rank never entered which collective. - Job-level metrics: step time per rank (p50, p99, max), data-wait time, tokens/s, MFU, loss, gradient norm, checkpoint durations, restarts.
- Straggler detection: compare per-rank compute time inside each step. A rank that is consistently slower (thermal throttling, a degraded NVLink, a bad NIC, a noisy neighbour on the host) stretches every collective. Flag ranks whose time exceeds the median by a threshold for k consecutive windows, then drain and replace the node.
In synchronous training a hang and a crash differ in the way that matters: a crash is loud, while a hang is silent. One rank stuck in a collective leaves 16k GPUs waiting until a timeout fires, which can take minutes to tens of minutes. Hang detection (step-progress heartbeats with a short deadline, plus collective tracing to find the culprit) is often worth more goodput than shaving checkpoint time.
- Kueue documentation: ClusterQueue, cohorts, borrowing, preemption, fair sharing, topology-aware scheduling.
- Slurm topology guide and preemption guide.
- MosaicML Streaming: deterministic, resumable shard streaming from object storage.
- Ray Data: streaming CPU+GPU pipelines.
- Fire-Flyer AI-HPC (DeepSeek): a cost-focused cluster design including the 3FS file system.
- dcgm-exporter: the usual starting point for GPU fleet metrics.
Part 2 · Case studies
Each case study follows the same arc: clarifying questions, scale estimates, an architecture diagram, a component walkthrough, deep dives, failure modes and trade-offs. In a real 45–60 minute interview you'll cover perhaps half of what is written here. Pick the deep dives the interviewer seems to care about.
Case 1 · Web-scale pre-training data pipeline
Prompt: "Design the pipeline that turns Common Crawl into a 15–30T-token pre-training corpus, and that the team can re-run when the filters change."
Clarifying questions
- Target size and quality bar: 15T tokens of the best data, or a larger pool to choose from later? English only or multilingual? Is code a separate pipeline?
- How often will filters change? (Constantly, during research. Design for cheap reprocessing.)
- Compliance: robots.txt and opt-out handling, PII policy, licensing, takedown requests that must propagate to future training sets.
- Compute: an existing Spark or Ray cluster? GPUs available for classifier inference? Budget?
Scale estimates
A single recent Common Crawl snapshot (September 2026) contains about 2.17 billion pages and 361 TiB of uncompressed WARC data Common Crawl 2026. Assumptions for the whole job:
- About 100 snapshots → ~200B pages and ~35 PiB raw WARC (less if you start from compressed WARC on a nearby cloud region).
- HTML-to-text extraction at ~5–20 ms of CPU per page → 200B × 10 ms = 2×10⁹ core-seconds ≈ 560k core-hours. At ~$0.01–0.03 per core-hour on spot, that is roughly $5–17k. Extraction is the largest CPU stage, but it's cheap compared with GPUs.
- After language ID, heuristic filters and dedup, maybe 10–20% of extracted text survives.
- Model-based quality scoring (a small classifier, or embeddings plus a linear head as in FineWeb-Edu) on ~20B candidate docs. A fastText classifier runs on CPU. A transformer scorer at ~5–10k docs/s per GPU needs 20B ÷ 7.5k ≈ 2.7M GPU-seconds ≈ ~750 GPU-hours.
Component walkthrough
- Ingestion and URL filtering. List the WARC files of each snapshot and partition the work by file, so each task processes one WARC file of about 1 GB. Apply URL blocklists (adult content, spam domains), robots and opt-out lists, and URL-level exact dedup before extraction. It is the cheapest filter and saves CPU downstream.
- Extraction. HTML to main-content text. RefinedWeb used trafilatura Penedo+ 2023. DCLM chose resiliparse, reporting it about 8× faster with comparable downstream results Li+ 2024.
- Language ID. A fastText classifier (lid.176, 176 languages) fastText runs at CPU speed. Store the language and its confidence as columns, and threshold later.
- Heuristic filters. Record each rule's result as a boolean or score column. Don't drop rows yet.
- Dedup. Exact (normalized-content hash), fuzzy (MinHash-LSH), and optionally line-level or substring-level. See the deep dive below.
- Model-based annotation. Quality classifiers (DCLM's fastText classifier, FineWeb-Edu's educational-value scorer trained on Llama-3-70B annotations Penedo+ 2024), toxicity classifiers, PII span detectors, and domain or topic labels used for mixing. These are GPU-friendly, so run them as Ray Data stages with actor pools of GPU workers.
- Selection and transformation. Apply thresholds, redact PII spans (emails, phone numbers, IP addresses) with placeholder tokens, and decontaminate: remove documents with long n-gram overlap with benchmark test sets (see A8).
- Tokenize and shard. Write token-ID arrays with a document index, sharded at hundreds of MB, plus a manifest of (shard → source table snapshot, row ranges, token counts per bucket).
Deep dive: MinHash-LSH dedup at trillions of tokens
Algorithm recap. Shingle each document into word n-grams (FineWeb uses 5-grams). Compute \(k\) MinHash values. For two documents with Jaccard similarity \(s\), each MinHash matches with probability \(s\). Split the \(k\) values into \(b\) bands of \(r\) rows. Two documents become a candidate pair if all \(r\) values match in at least one band, which happens with probability \(1 - (1 - s^r)^b\). FineWeb uses 112 hashes as 14 bands of 8, targeting documents at least 75% similar Penedo+ 2024. Check the S-curve: \(s = 0.5\) gives about 5%, \(s = 0.75\) about 77%, \(s = 0.8\) about 92%, \(s = 0.9\) about 99.96%. The threshold sits near \((1/b)^{1/r} \approx 0.72\). The standard reference for banding is chapter 3 of Mining of Massive Datasets.
Distributing it. This is four jobs, and the second one is a shuffle:
- Signatures (map, embarrassingly parallel): per document, compute the \(b\) band hashes, each a 64-bit hash of the \(r\) MinHash values. Emit \(b\) records of
(band_id, band_hash, doc_id). - Bucket (shuffle and group-by): partition by
(band_id, band_hash). Within each bucket, every document shares a band with every other, so emit edges. Emit a star (each member → the bucket's minimum doc_id), not all pairs, to avoid quadratic blow-up on huge buckets of boilerplate pages. - Connected components: union the edges into clusters. At moderate scale, collect the edges (much smaller than the documents) onto one large machine and run union-find. At larger scale, use an iterative distributed connected-components algorithm (label propagation or hash-to-min style) that converges in a logarithmic number of rounds.
- Filter (join): keep one representative per cluster (for example the most recent crawl, or the highest quality score), and anti-join the rest out of the table.
Hugging Face's datatrove implements MinHash dedup as separate signature, bucket, cluster and filter stages that can run on a local machine, Slurm or Ray datatrove.
Shuffle cost. Take 100T tokens at ~1,000 tokens per doc, so about \(10^{11}\) docs. Bucket records: 14 bands × 10¹¹ docs × ~16 bytes (8 B hash + 8 B doc_id, band_id implied by partition) ≈ 22 TB shuffled. That is a big but routine Spark shuffle, roughly an hour or two on a few hundred nodes with local SSD for spill (rough). Skew is the real danger: a boilerplate template ("404 page not found", cookie walls) can produce a single bucket with hundreds of millions of members. Handle it with the star emission above, plus salting or capping very large buckets, which are themselves a strong signal of spam.
Global dedup isn't automatically better. FineWeb found that deduplicating each snapshot on its own beat deduplicating across all 96 snapshots. Global dedup removed proportionally more of the good recent data and kept low-quality pages that happened to be unique Penedo+ 2024. Llama 3, by contrast, reports global MinHash dedup across the whole dataset, plus line-level dedup that removes lines appearing more than 6 times per 30M-document bucket Llama Team 2024. The interview point: dedup scope and threshold are hyperparameters to ablate, not fixed facts. The pipeline must make them cheap to change.
Exact and substring dedup. Exact content hashes are a cheap group-by. For removing repeated spans such as boilerplate paragraphs and licence text, Lee+ used suffix arrays to find exact substrings repeated across documents Lee+ 2021. RefinedWeb combined MinHash with exact substring dedup Penedo+ 2023. DCLM used a Bloom-filter approach (BFF) for scalable paragraph- and document-level dedup Li+ 2024.
Deep dive: reprocessing when filters change
- Annotate, don't delete. Every filter writes a score column to a wide table keyed by doc_id. "Change the quality threshold from 0.6 to 0.7" is then a cheap selection job, not a re-crawl.
- Version every stage. Output tables are keyed by
(stage, stage_code_version, config_hash, input_snapshot_id). The orchestrator (Airflow, Dagster, Flyte or in-house) recomputes only stages whose key changed and their descendants. It works like a build system with content-addressed caching. - Dedup is not incremental-friendly. Adding a new snapshot changes clusters. Options: dedup new data against the existing corpus's signature index (only new docs can be removed), or periodically recompute globally. Per-snapshot dedup (the FineWeb choice) makes adding a snapshot independent of the others.
- Takedowns and opt-outs must propagate. Keep a doc_id → shard index so you can produce a new dataset version without the removed docs, and record which training runs used the old version.
Failure modes
- Poison inputs: a malformed WARC record or a 2 GB HTML page crashes a worker. Use per-record try/except, size caps, a quarantine table, and task-level retries with a bounded attempt count.
- Partial outputs after task failure: write to a temp path and commit atomically, through a table-format commit or a manifest swap. Re-running a task must be idempotent.
- Silent filter bugs, such as a language-ID threshold that drops all of Hindi. Track per-stage survival rates by language, domain and crawl, and alert on shifts.
Interviewers often push on "how do you know your filters improve the model?" The infra answer: the pipeline exists to feed an ablation loop. Train small proxy models (around 1B parameters on tens of billions of tokens) on alternative dataset versions and compare downstream evals. FineWeb and DCLM both work this way. So the pipeline must be able to produce many dataset variants cheaply. That requirement is what justifies the annotate-don't-delete design and stage-level caching.
- FineWeb: the most detailed public account of a web pipeline, including dedup ablations.
- DataComp-LM (DCLM): a benchmark for data curation, with resiliparse extraction, BFF dedup and a fastText quality filter.
- datatrove and NeMo Curator: open-source pipeline frameworks (CPU, and GPU-accelerated).
- Mining of Massive Datasets, ch. 3: MinHash and LSH banding from first principles.
- A2 · Pre-training: why each filtering stage matters for model quality.
Case 2 · Orchestrating a 16k-GPU training run with fault tolerance
Prompt: "You're running a 16,384-GPU pre-training job for three months. Design the system that keeps goodput above 90%."
Clarifying questions
- Model size and parallelism layout (TP × PP × DP, plus EP for MoE). This decides checkpoint size and what one failure takes down.
- Hardware: GPU generation, nodes per scalable unit, network (InfiniBand or RoCE), storage system and its bandwidth.
- Do we own the cluster, or is it a cloud reservation? Are spare nodes available, and how fast can a node be repaired or replaced?
Scale estimates
- Failure rate. Llama 3 data gives about one unexpected interruption every 3.1 h at 16k GPUs (Part 0). Meta's research-cluster study observed a 7.9-hour MTTF for 1,024-GPU jobs and projected about 1.8 hours for 16,384-GPU jobs Kokolis+ 2024. Plan for M ≈ 2–3 h, so 8–12 interruptions a day.
- Checkpoint size. Take a 405B dense model at ~14 bytes per parameter of training state: \(405\text{B} \times 14 \approx 5.7\) TB. Per GPU, if state is fully sharded: 5.7 TB ÷ 16,384 ≈ 350 MB. Llama 3 reports checkpoint sizes of 1 MB to 4 GB per GPU, depending on rank Llama Team 2024.
- Checkpoint write time. Synchronous to a storage system sustaining 2 TB/s (Llama 3's Tectonic figure): 5.7 TB ÷ 2 TB/s ≈ 3 s if perfectly parallel. In practice, with metadata overhead, stragglers and contention, plan for 30–120 s. Async: training blocks only for the GPU→host copy, which is 350 MB per GPU at ~20 GB/s effective PCIe, about 20 ms plus synchronization, so a few seconds visible.
- Restart time R. Detect (timeout or heartbeat), drain the bad node, allocate a spare, launch containers, NCCL init across 16k ranks, load the checkpoint, warm up. Naively 15–30 min. Optimized: ~3–5 min.
Goodput math. Use the Young/Daly interval \(\tau^* = \sqrt{2CM}\) Young 1974 Daly 2006 and the first-order loss model:
$$\text{overhead} \approx \frac{C}{\tau} + \frac{\tau/2 + R}{M}$$| Scenario | C (blocking) | R | τ* = √(2CM), M = 180 min | Overhead | ≈ Availability |
|---|---|---|---|---|---|
| Naive: sync checkpoint, slow restart | 2 min | 25 min | 27 min | 7.4% + 21.4% ≈ 29% | ~71% |
| Sync checkpoint to fast storage, decent restart | 0.5 min | 12 min | 13 min | 3.8% + 10.3% ≈ 14% | ~86% |
| Async + in-memory checkpoint, hot spares | 0.05 min | 4 min | 4.2 min | 1.2% + 3.4% ≈ 4.6% | ~95% |
The table shows that once checkpoints are async, restart time dominates. Every minute of R costs about 0.55% of availability at M = 3 h. Meta's research-cluster analysis reaches a similar conclusion: for 12k-GPU jobs to reach an effective training time ratio of 0.9, checkpoint overhead must fall to about 10 s or the failure rate must fall substantially Kokolis+ 2024.
Component walkthrough
- Job controller. Holds the desired state (run ID, config, node count, latest committed checkpoint) in a strongly consistent store and drives a state machine. Restarts are automatic. Humans get paged only for repeated failures or loss anomalies. Llama 3 reports needing significant manual intervention only three times in the 54-day window Llama Team 2024.
- Node health service. Before a node joins a job it passes burn-in: GPU diagnostics, NVLink and NIC loopback, and small NCCL all-reduce tests against neighbours. MegaScale describes this kind of self-check suite (intra-host loopback and RDMA NIC tests, NCCL tests) plus a training daemon that heartbeats status and logs to a driver Jiang+ 2024. While running, the service watches DCGM, XID errors and ECC counters. A node that fails goes to quarantine, gets triaged, and then goes back to the pool or out for RMA.
- Hot spares. Size the pool with Little's law: spares needed ≈ failure rate × repair time. With about 8 node failures a day and 2-day repair turnaround, roughly 16 nodes are out at any time. Keep about 1–2% spare capacity (for example 32 nodes for 2,048 active), already burned in, with images pre-pulled and placed so a swap keeps the topology tight.
- Hang and failure detection. Each rank reports a step heartbeat. If global step progress stalls for longer than k × the normal step time, declare a hang, dump flight-recorder traces from all ranks, and find the rank or ranks that never entered the pending collective. Llama 3 describes using PyTorch's NCCL flight recorder this way Llama Team 2024.
- Checkpointing. Use sharded checkpoints, where each rank writes its own shard and global metadata allows load-time resharding onto a different topology PyTorch DCP. Make them async: stage to host RAM and persist in the background. PyTorch reported cutting visible checkpoint downtime for a 7B model from 148.8 s to 6.3 s this way PyTorch 2024. Add in-memory redundancy: Gemini kept redundant in-memory copies of model state, recovered from an intact replica, and reports that goodput for its largest job rose from 85% to 97% Gemini Team 2023. ByteDance's ByteCheckpoint is a production system along the same lines Wan+ 2024.
- Commit protocol. A checkpoint counts only when every shard is durable and a manifest (step, shard list, checksums, data-loader state, RNG state, config hash) is atomically written. The controller restarts from the latest committed manifest.
Deep dive: stragglers
In synchronous training every collective waits for the slowest rank. MegaScale records per-rank execution time of critical code segments with CUDA events and visualizes them as heatmaps. It found about 0.5% of machines were substantially slower Jiang+ 2024. A ByteDance trace study found that stragglers slowed 42.5% of jobs by at least 10% and wasted about 10% of GPU-hours. Most of the straggling did not come from individual bad machines. Pipeline-stage imbalance, sequence-length variance across micro-batches, and uncoordinated Python garbage-collection pauses were bigger causes Lin+ 2025. So straggler handling has two halves:
- Hardware stragglers: detect a rank persistently above the median (for example 1.1× for 10 consecutive windows), cross-check with DCGM (clocks, temperature, PCIe or NVLink errors), then drain and swap at the next checkpoint boundary.
- Software stragglers: balance pipeline stages, pack sequences to even token counts per micro-batch, and synchronize or disable automatic GC (run
gc.collect()on all ranks at the same step).
Deep dive: elastic and fault-tolerant training
Checkpoint-and-restart treats the whole job as one failure domain. The alternative is to keep training through a failure. torchft provides per-step fault tolerance: data-parallel replica groups are independently recoverable, a Lighthouse coordinator tracks health by heartbeat, and a failed group drops out of the cross-group gradient all-reduce and later rejoins by fetching weights from a healthy peer torchft. In a PyTorch demonstration on 300 L40S GPUs (30 replica groups), training continued with checkpointing turned off while failures were injected every 60 s (about 82% step efficiency) and even every ~15 s, and the model still converged PyTorch 2025. Trade-offs to state:
- Dropping a replica changes the effective global batch for that step. Either rescale the learning rate or gradient, or accept the noise.
- Model-parallel groups (TP, PP) can't lose a member. The unit of failure is a whole replica group, which for a 405B model may be dozens of nodes.
- Semi-synchronous methods (LocalSGD, DiLoCo-style) go further by syncing only every H steps, which tolerates slow or failed islands and cross-datacenter links.
Fault-tolerant and elastic training is an active area (torchft, ByteDance's ByteRobust reporting 97% effective training time over three months on 9,600 GPUs ByteRobust 2025, in-memory checkpoint systems, semi-synchronous optimizers). Check the current PyTorch, Megatron and lab reports for production practice.
Other failure modes
- Silent data corruption (SDC). A GPU computes wrong numbers without raising an error. Gemini reported expecting SDC events to affect training every week or two at its scale, and used deterministic replay plus proactive scanners on idle machines and hot standbys to catch them Gemini Team 2023. Symptoms: loss or grad-norm spikes, NaNs, or DP replicas whose gradients disagree.
- Loss spikes or divergence. Not a hardware fault. Roll back to the last good checkpoint, skip the offending data window (the deterministic loader makes this a config change), and maybe lower the learning rate. Keep the rollback path automated but human-approved.
- Power swings. Llama 3 notes that the whole job pausing for checkpoints or collectives causes data-centre power swings on the order of tens of megawatts Llama Team 2024. At this scale the power grid and cooling are part of your system.
Run dashboard and alerting
| Panel | Signal | Alert |
|---|---|---|
| Progress | Tokens/s, step time, MFU, tokens trained vs plan | MFU drop > 10% sustained for 15 min |
| Health of the model | Loss (train and held-out), grad norm, update/weight ratio, NaN count | Loss spike > k σ, NaN, grad-norm explosion |
| Goodput breakdown | Time in train, checkpoint, restart, hang, idle-waiting-for-nodes | Daily goodput below target |
| Stragglers | Per-rank compute time heatmap, p99/p50 ratio | Same node above threshold for N windows |
| Fleet | Healthy, suspect, quarantined and spare node counts. XID and ECC rates | Spare pool below minimum |
| Checkpoints | Last committed step, save and upload durations | No committed checkpoint in 2× expected interval |
The headline trap: candidates talk about checkpoint write speed, but at scale goodput is lost mainly to detection plus restart. Show the table above, then say where the remaining minutes go: timeouts that fire late, container pulls, NCCL bootstrap over thousands of ranks, checkpoint loading from cold storage. Then propose fixes for each: shorter hang deadlines with heartbeats, pre-pulled images on spares, a persistent NCCL or rendezvous service, restore from peer RAM.
- Llama 3 report §3.3.4: interruption taxonomy and tooling.
- MegaScale §4–5: robust training daemon, self-checks, two-stage checkpointing, straggler heatmaps.
- Revisiting Reliability in Large-Scale ML Research Clusters: MTTF by job size and the effective-training-time model.
- Understanding Stragglers in Large Model Training: where slowness really comes from.
- torchft and the fault-tolerant Llama blog post.
- PyTorch Flight Recorder tutorial: debugging stuck collectives.
Case 3 · Multi-tenant GPU cluster platform for a research org
Prompt: "We have 4,096 GPUs and 300 researchers across 12 teams. Design the platform: how jobs get submitted, scheduled and charged. Today people complain about both long queues and idle GPUs."
Clarifying questions
- Workload mix by GPU-hours and by job count. The Acme study of two ~2,300-GPU LLM clusters found pre-training was 3.2% of jobs but 94% of GPU time in one cluster, while evaluation jobs were 92.9% of jobs but 0.8% of resources, and the median job ran for about 2 minutes Hu+ 2024. Expect a similar split: a few huge jobs and a flood of tiny ones.
- Are allocations contractual (team X funded 20% of the cluster) or purely priority-based?
- Interactive needs: notebooks, debugging on 1–8 GPUs, with what latency expectations?
Scale estimates
- 4,096 GPUs = 512 nodes of 8. At ~$2.5 per GPU-hour equivalent cost, the cluster is worth about $10k per hour, or ~$90M per year. Each 1% of utilization is worth about $0.9M a year, which is the business case for platform work.
- Job arrival: thousands of jobs a day, mostly small.
- Big jobs: 1–4 runs of 512–2,048 GPUs that need topology-tight placement.
Component walkthrough
- Quota model. Each team gets a nominal quota (GPUs by type). Teams in a cohort can borrow each other's unused quota, and the lender can reclaim it through preemption. This is Kueue's ClusterQueue and cohort model, with nominal quotas plus borrowing and lending limits Kueue. On Slurm the equivalent is QOS limits plus fair-share and partition priorities. Quotas live in Git and change through code review.
- Priorities. Three to five classes are enough: reserved production > team-guaranteed > team-borrowed > interactive > scavenger. Kubernetes PriorityClasses drive pod-level preemption Kubernetes docs. Kueue preempts at workload level, so it evicts whole gangs rather than random pods.
- Gang admission. A workload is admitted only if its full resource request fits. Volcano expresses this with a PodGroup minMember Volcano docs. Kueue admits whole workloads against quota, and JobSet describes multi-role jobs.
- Topology. Label nodes by NVLink domain, leaf switch and block. Large jobs request "same block". Kueue's topology-aware scheduling Kueue TAS and Slurm's tree and block topology plugins Slurm docs implement this.
- Checkpoint-aware preemption. Preemption sends SIGTERM with a grace period (for example 2–5 minutes) during which the job writes an emergency checkpoint. Jobs declare whether they are preemptible and how long they need.
- Interactive pool. A separate partition with GPU-granular allocation (1–8 GPUs), short time limits, and an idle reaper: if SM activity stays near zero for 30–60 minutes, warn, then reclaim. Interactive sessions are the biggest source of "allocated but idle".
- Scavenger tier. Fill every gap with preemptible work such as eval sweeps, hyperparameter searches and data-processing GPU stages. It is evicted first and must checkpoint frequently.
Deep dive: fragmentation
Scenario: 100 nodes are free, but a 64-node job can't start, because the free nodes are spread across 8 blocks and the job requires topology-tight placement. Or 400 GPUs are free, but in 1–3-GPU slivers on partially used nodes, so no 8-GPU-per-node job fits. Remedies:
- Segregate by shape: small jobs go to a bin-packing partition, and multi-node jobs get whole nodes elsewhere. Never put a 1-GPU job on an otherwise empty node in the large-job pool.
- Packing policy: best-fit (place small work where it leaves the fewest broken blocks) rather than spreading.
- Reservations and backfill: when a big job is at the head of the queue, reserve nodes as they free up, and backfill only jobs that will finish before the reservation's start time.
- Defragmentation: periodically preempt and requeue small preemptible jobs so they consolidate.
- Measure it: report the largest gang size that could start now, alongside free GPUs.
Deep dive: chargeback
Chargeback (or showback) bills teams for allocated GPU-hours at a rate that depends on tier: reserved costs full price, borrowed costs less, scavenger is nearly free. This creates incentives to release idle allocations and use preemptible tiers. Publish per-team dashboards of allocated vs active GPU-hours. Making idleness visible often recovers more capacity than any scheduler change.
Failure modes and trade-offs
- Preemption storms: a big reserved job starts and evicts 300 small jobs at once, all of which requeue and thrash. Rate-limit preemptions and prefer victims that checkpointed recently.
- Trade-off: strict guarantees (static partitions) are predictable but waste capacity. Full sharing maximizes utilization but makes large-run start times unpredictable. Most orgs use guarantees with borrowing plus a reserved block for flagship runs.
When asked "utilization is 60%, fix it", first separate allocation from activity. Low allocation means a scheduling problem (fragmentation, quota walls, slow startup). High allocation with low activity means a workload problem (idle notebooks, data stalls, hangs). The fixes are completely different, and saying so early signals experience.
- Characterization of LLM Development in the Datacenter (Acme, NSDI'24): real workload mix, failure and utilization data from two LLM clusters.
- Kueue concepts: ClusterQueue, cohort, preemption, fair sharing, TAS.
- Volcano PodGroup: gang semantics on Kubernetes.
- Slurm Fair Tree and preemption.
- Building Meta's GenAI Infrastructure: two 24k-GPU clusters (RoCE and InfiniBand) and the storage behind them.
Case 4 · RL training infrastructure at scale
Prompt: "Design the infrastructure for RL post-training of a 70B model on math, coding and agentic software-engineering tasks, with verifiable rewards and up to thousands of concurrent sandboxed environments." For the algorithms (GRPO, PPO, staleness corrections), see A5. This section is about the systems.
Clarifying questions
- Algorithm family: group-based (GRPO-style, no critic) or PPO with a value model? A critic adds a model to train and to place.
- Rollout lengths: single-turn math (2–16K tokens) or multi-turn agentic episodes (dozens of tool calls, minutes of wall-clock time each)?
- Reward sources: programmatic verifiers (unit tests, answer checkers), reward models, LLM judges? What are their latencies?
- How much off-policyness is acceptable? Will the research team tolerate asynchronous training?
- GPU budget, and whether rollout and training can share GPUs (colocated) or must be separate pools.
Scale estimates
Assume each step uses 512 prompts × 16 samples = 8,192 rollouts with an average of 8K generated tokens. That is about 67M generated tokens per step.
- Generation. A 70B model in FP8 on one 8×H100 node at ~8K context: KV cache per token is about 2 × 80 layers × 8 KV heads × 128 dims × 2 bytes ≈ 330 KB (bf16 KV). Memory left for KV after weights is ~550 GB, so about 1.6M tokens of KV, or ~200 concurrent 8K sequences. Each decode step reads weights (70 GB) plus KV (up to ~500 GB) at an aggregate ~27 TB/s, so ~20 ms per step for ~200 tokens. That is ~5–10K tokens/s per node. 67M tokens ÷ 7.5K tokens/s ≈ 9,000 node-seconds: ~70 s on 128 rollout nodes, if the work were perfectly balanced.
- Training. 6 × 70B × 67M ≈ 2.8×10¹⁹ FLOPs, plus about 2N per token each for the reference-policy and old-policy log-prob passes. Call it 10N per token ≈ 4.7×10¹⁹ FLOPs. At ~400 TFLOPS effective per GPU: ~460 s on 256 GPUs, ~115 s on 1,024.
- Takeaway. Per GPU, generation and training throughput are in the same ballpark at 8K context. Generation dominates wall-clock time in practice because of the long tail (a few 32K-token rollouts hold up the whole batch), because KV pressure at long contexts shrinks the batch, and because of environment latency in agentic tasks. Size the pools from measured numbers, not folklore.
- Weight sync. 70B in bf16 = 140 GB per step. A pipelined broadcast over 400 GB/s per-node InfiniBand takes a few seconds to reach 128 nodes. Pulling from object storage would be 140 GB × 128 = 18 TB of fan-out per step, so don't do that.
- Sandboxes. Say 25% of tasks are agentic, with ~5 minutes per episode, and the async pipeline keeps ~4,000 episodes in flight. At 1–2 vCPU and 2–4 GB RAM each, that is ~6,000 vCPU, about 50–100 CPU nodes.
Component walkthrough
- Placement. Colocated: the same GPUs alternate between generation and training. veRL's HybridFlow reshards actor weights between training and generation layouts with its 3D-HybridEngine Sheng+ 2024, and vLLM's sleep mode can offload weights and free the KV cache between phases vLLM docs. Disaggregated: separate pools that overlap in time, sized independently. OpenRLHF is built on Ray, vLLM and DeepSpeed and places roles as Ray actors Hu+ 2024. Kimi k1.5 reports switching from training to inference in under a minute, and about ten seconds in the other direction, in its hybrid deployment Kimi Team 2025.
- Synchrony. Asynchronous designs decouple generation from training. AReaL is fully asynchronous with a maximum-staleness hyperparameter η and a staleness-aware PPO variant, and reports up to 2.77× speedup over synchronous training on the same GPUs Fu+ 2025. PipelineRL updates weights in flight, mid-generation, and reports about 2× faster learning on 128 H100s Piché+ 2025. The infra consequence: every trajectory must carry the policy version (or versions) that generated it, plus the rollout log-probs, so the trainer can filter by staleness and apply importance correction.
- Long-tail handling. Partial rollouts cap the generation token budget per iteration and carry unfinished sequences into the next one through a replay buffer Kimi Team 2025. Kimi K2 reports using partial rollouts for long-horizon agentic tasks too Kimi Team 2025b. Alternatives are over-sampling and dropping stragglers (this biases against long answers) and hard length caps.
- Weight sync. Reshard from the trainer's FSDP or TP/PP layout to the engine's TP layout, then broadcast in buckets. Use NCCL across nodes, CUDA IPC when colocated, or ship LoRA deltas only. Kimi k1.5 moves weights between nodes with Mooncake over RDMA Kimi Team 2025. Version every push so engines can report which version produced each token.
- Reward service. A separate autoscaled service with a uniform API:
score(trajectory) → reward, metadata. Verifiers that run code reuse the sandbox pool. Reward models and LLM judges run as inference replicas.
Deep dive: an environment service for thousands of sandboxes
Public reports give a sense of scale. Kimi K2 describes a Kubernetes-based sandbox system that supports over 10,000 concurrent sandbox instances Kimi Team 2025b. Kimi k1.5's code sandbox replaced Docker with crun, reused pre-created cgroups, and used an overlay filesystem with a tmpfs upper layer. Container startup dropped from 0.12 s to 0.04 s, and a 16-core machine could start 120 sandboxes per second instead of 27 Kimi Team 2025. Design points:
- Isolation level. Containers (fast, shared kernel), gVisor (user-space kernel, stronger isolation) gVisor, or Firecracker microVMs (hardware isolation, under 125 ms boot, under 5 MiB overhead, up to 150 microVMs per second per host) Firecracker. Policy-generated code is untrusted and, under RL pressure, adversarial: it will find reward hacks, and it may find escapes. Disable the network by default, cap resources, use read-only base images and per-episode ephemeral writable layers.
- Image management is the hidden cost. SWE-Gym ships pre-built Docker images per task instance totalling 6 TB Pan+ 2024. SWE-smith instead uses one environment per repository, 295 GB total, and estimates 50–150 TB for the per-task approach at its 50k-instance scale Yang+ 2025. Share base layers, pre-pull onto nodes, use lazy-loading snapshotters, and schedule episodes onto nodes that already have the image (image-affinity scheduling).
- Warm pools via Little's law. With 4,000 concurrent episodes of ~300 s mean duration, sandbox demand is λ ≈ 4,000 ÷ 300 ≈ 13 per second. If a cold start (pull, create, repo setup) takes 30 s, then to hide it you need at least λ × 30 ≈ 400 warm sandboxes per active image mix. Add burst headroom (×2) for the start of each step in a synchronous setup. Snapshot and restore of a fully set-up environment (repo checked out, dependencies installed) converts a 30 s setup into sub-second restore.
- Lifecycle and leaks. Every sandbox gets a TTL and a heartbeat from its episode driver. A reaper kills orphans. Without that, crashed drivers leak thousands of sandboxes within hours.
Failure modes
- Engine and trainer log-prob mismatch (different kernels and precision) makes "on-policy" data subtly off-policy. Log both, monitor the ratio distribution, and correct with truncated importance sampling (see A5).
- Reward-service outages masquerading as learning signal: a verifier fleet returning errors that are scored as 0. Alert on reward error rates per environment and quarantine affected batches.
- Reward hacking detected only by humans reading rollouts. Sample trajectories into a review UI continuously, and alarm on sudden reward jumps paired with length or behaviour shifts.
A strong answer makes the GPU-idle accounting explicit. In a synchronous colocated design, idle time is the long tail of each generation phase. In a disaggregated synchronous design, it's whichever pool finishes first. In an async design, idle time disappears, but you pay in staleness. Then say which you'd pick: async with a small η plus partial rollouts for long-CoT tasks, and colocation only if GPUs are scarce. Also mention that the environment fleet is CPU-bound, autoscaled separately, and often the true bottleneck for agentic RL.
- A5 · RL for LLMs §11: the conceptual design axes (placement, synchrony, weight timing).
- HybridFlow / veRL and its repo verl; OpenRLHF; slime (Megatron + SGLang); NeMo RL.
- AReaL and PipelineRL: asynchronous and in-flight-update designs.
- Kimi k1.5 §2.6: partial rollouts, hybrid deployment, the code sandbox.
- SWE-smith: environment construction at scale and why per-repo images matter.
Case 5 · Multi-model inference serving platform
Prompt: "Design an internal platform that serves ~30 LLMs (8B to 400B+ MoE) plus hundreds of fine-tuned LoRA variants, for both interactive products and internal batch users, with SLO tiers." Engine internals (KV cache, batching, disaggregation) are in A6. Here we design the platform around the engines.
Clarifying questions
- Traffic per model: peak QPS, input and output length distributions, diurnal shape. Is traffic concentrated in a few models (it usually is)?
- SLOs per tier: interactive (time to first token p99 < 1 s, inter-token latency < 50 ms), standard, batch (hours).
- Multi-region? Data-residency constraints? Are GPU pools fixed or elastic (cloud)?
Scale estimates
- Hot model: 3,000 requests/s at peak, 2,000 input and 300 output tokens on average. That gives 900K output tokens/s and 6M prefill tokens/s.
- A 70B-class model on an 8×H100 node under interactive SLOs might sustain ~5–15K output tokens/s with prefill mixed in. 900K ÷ 10K ≈ 90 nodes (720 GPUs) for one hot model, before headroom.
- Long-tail models: 25 models at under 10 requests/s each. Dedicating even one node to each wastes capacity. They need consolidation (multi-model per node, scale to zero, LoRA).
- Cold start. A 70B model in bf16 is 140 GB. From object storage at ~2–5 GB/s that takes 30–70 s. From local NVMe at ~10–20 GB/s it takes 7–14 s. Add engine init, CUDA graph capture and warmup, and the end-to-end cold start is ~1–5 minutes. Autoscaling can't absorb sub-minute spikes, so you need headroom or predictive scaling.
Component walkthrough
- Registry to deployment. A model version is an immutable artifact (safetensors weights, tokenizer, chat template, engine config, eval report). A deployment spec names the version, hardware, parallelism, quantization and autoscaling policy. Aliases ("chat-default") are mutable pointers moved by the canary process. Safetensors avoids pickle's arbitrary-code risk and supports lazy loading safetensors. KServe and NVIDIA Dynamo are open-source examples of the deployment layer on Kubernetes KServe Dynamo.
- Routing. Plain round-robin wastes prefix cache. KV- or prefix-aware routing sends a request to the replica most likely to hold its prefix (system prompt, conversation history), weighted against load. Dynamo's router routes on worker load and KV-cache overlap Dynamo, and llm-d offers prefix-cache and load-aware balancing llm-d. The Kubernetes Gateway API Inference Extension standardizes an "endpoint picker" for model-server pools Gateway API Inference Extension. Use consistent hashing on a prefix hash, with bounded load, so one hot prefix doesn't overload a replica.
- Prefill/decode disaggregation for hot models with long prompts: separate prefill and decode pools sized independently, with KV transferred between them. See DistServe Zhong+ 2024 and Mooncake Qin+ 2024, and A6 for when it pays off.
- LoRA multi-adapter serving. Hundreds of fine-tunes of one base model share the base weights. Batching requests for different adapters together with custom kernels is the idea behind Punica Chen+ 2023 and S-LoRA ("thousands of concurrent adapters") Sheng+ 2023. vLLM exposes this with
--enable-lora,--max-lorasand--max-cpu-lorasvLLM docs. Platform jobs: an adapter cache hierarchy (GPU slots → host RAM → object store), adapter-affinity routing, and admission when a request's adapter isn't resident. - Autoscaling signals. CPU utilization is meaningless here. Use queue depth per priority, KV-cache utilization (near 100% means requests are about to queue or get preempted), and TTFT and ITL against SLO. Scale up early because of the long cold start. Scale down slowly (hysteresis of tens of minutes). Pre-warm on the daily traffic pattern.
- Cold start reduction. Cache weights on local NVMe across restarts. Stream tensors straight to GPU with concurrent reads: vLLM supports the Run:ai Model Streamer from S3, GCS and Azure vLLM docs. Copy weights GPU-to-GPU from a running replica (Dynamo describes streaming weights via NIXL/NVLink Dynamo). Keep CUDA graphs and compiled artifacts cached.
- SLO tiers. Map tiers to priority queues with preemption: interactive requests can preempt batch-tier sequences in the engine scheduler, or the router can shed batch traffic first. Batch-tier traffic soaks up idle capacity at night and is priced lower. Public batch APIs offer about 50% discounts for 24-hour turnaround OpenAI docs Anthropic docs.
- Canary rollouts. A new version moves through 1% → 5% → 25% → 100% of traffic, with each stage gated on latency SLOs, error rates, and quality signals (online eval sampling, judge-scored diffs on mirrored traffic). The previous version stays warm so you can roll back instantly.
Failure modes
- KV-cache thrash: too many long requests cause engine preemptions and recomputation, and throughput collapses. Use admission control on estimated KV footprint (prompt length + max tokens).
- Silent quality regression from a quantization or engine upgrade: latency is fine but answers got worse. Run eval gates before deploy and shadow-compare during canary.
The serving control-plane ecosystem (Dynamo, llm-d, AIBrix AIBrix, KServe, the Gateway API Inference Extension) changed quickly through 2025–2026. Check the current feature sets before naming specific capabilities in an interview.
- A6 · Inference & Serving: engine-level concepts this platform relies on.
- NVIDIA Dynamo and llm-d: open-source distributed serving with KV-aware routing and disaggregation.
- S-LoRA and Punica: multi-adapter serving.
- Mooncake: a KV-cache-centric production architecture.
- Gateway API Inference Extension: Kubernetes-native model routing.
Case 6 · Batch inference: embedding backfill for 1B documents
Prompt: "We're switching embedding models. Re-embed our 1B-document corpus and load it into the vector store, with no downtime for search. Minimize cost."
Clarifying questions
- Document length distribution and chunking policy. How many chunks per doc?
- Embedding model size and output dimension. Will we store reduced precision (fp16, int8, binary)?
- Deadline (a day, a week)? A longer deadline lets you use cheaper spot GPUs and smaller GPU types.
- Vector store: managed or self-hosted? Can it bulk-import from files, or only accept upserts?
Scale and cost estimates
- 1B docs × 1.5 chunks × 512 tokens ≈ 7.7×10¹¹ tokens.
- Model: a ~0.5B-parameter encoder. FLOPs ≈ 2 × 0.5B × 7.7×10¹¹ ≈ 7.7×10²⁰, plus attention overhead, call it 10²¹.
- On H100 at a realistic ~30% MFU (≈300 TFLOPS): 10²¹ ÷ 3×10¹⁴ ≈ 3.3×10⁶ GPU-s ≈ ~930 GPU-hours. At ~$2.5/h that is about $2.3k. On cheaper inference GPUs (L4/A10-class) you need more GPU-hours at a lower rate. Benchmark tokens per dollar, not tokens per second.
- Output: 1.5B vectors × 1,024 dims × 2 bytes (fp16) ≈ 3 TB, plus IDs and metadata.
- Vector store ingestion: if the GPU phase takes 10 hours, the upsert rate needed is 1.5B ÷ 36,000 s ≈ 42k vectors/s. Many vector databases can't sustain that while also serving queries and building HNSW graphs. Ingestion, not GPUs, is often the bottleneck.
- Contrast with LLM batch generation, for example summarizing the same 1B docs with an 8B model at 1,000 tokens in and 200 out. Prefill: 2 × 8B × 10¹² ≈ 1.6×10²² FLOPs ≈ 15k GPU-hours at 30% MFU. Decode: 2×10¹¹ tokens at ~2–4K tokens/s per GPU ≈ 14–28k GPU-hours. Total about 30–45k GPU-hours, roughly $75–110k, about 30–50× the embedding job. The same platform design applies.
Component walkthrough
- Pin a snapshot. Embed snapshot S of the corpus, not "the live table". Changes after S are handled by a catch-up job fed by change data capture (CDC).
- Sharding and leases. Shards sized to about 2–5 minutes of GPU work. That is small enough that a spot interruption (AWS gives a 2-minute warning AWS docs, and GCP's Spot VM shutdown period is best effort and up to 30 s GCP docs) loses little, and large enough that per-shard overhead stays negligible. Workers claim shards by taking a lease with an expiry, heartbeat it, and mark it complete with the output URI. Expired leases make shards claimable again. Cap attempts and send poison shards to a dead-letter state.
- Idempotent output. Output path = f(model_version, snapshot, shard_id). Write to a temp key, then rename or commit. A duplicated shard (a zombie worker finishing after its lease expired) overwrites identical content, which is harmless.
- Throughput engine. Sort or bucket by length to minimize padding, use large batches, and overlap CPU tokenization with GPU compute. Hugging Face's text-embeddings-inference does token-based dynamic batching TEI. Ray Data can run the whole read → tokenize → encode → write pipeline with CPU and GPU stages, and has built-in LLM batch processors on vLLM and SGLang Ray docs.
- Index build and cut-over. Build the new index offline from the Parquet files, using a bulk-import path if the store has one. At 1.5B vectors, full-precision HNSW in RAM is about 3 TB of vectors plus graph overhead, so consider IVF-PQ or disk-based graph indexes (see B2). Validate recall on a labelled query set. Run the catch-up delta. Flip the alias the search service reads. Keep the old collection for rollback.
Failure modes and trade-offs
- Spot capacity disappears: fall back to on-demand above a deadline-driven threshold. The planner knows the remaining work and the time left.
- Trade-off: online upserts into the live collection are simpler, but they slow live search and make rollback hard. Offline build plus alias flip costs double storage temporarily but is safe.
Interviewers like the moment you notice the GPU part costs a couple of thousand dollars and finishes in hours, while the vector-store load and validation take longer and carry more risk. Say it, then spend your time on the cut-over, the CDC catch-up and rollback.
- Ray Data: Working with LLMs: batch inference pipelines on vLLM or SGLang.
- text-embeddings-inference: a throughput-oriented embedding server.
- OpenAI Batch API and Anthropic batch processing: how hosted batch tiers are priced and limited.
- B2 · Retrieval: index types and their memory and recall trade-offs.
Case 7 · Evaluation platform for training checkpoints
Prompt: "Every checkpoint from our training runs should be evaluated on ~200 benchmarks, including sandboxed code execution and LLM-judge evals. Design the platform, including dashboards and regression alerts." For what to measure and statistical rigor, see A8.
Clarifying questions
- How many runs and checkpoints a day? How fast do results need to come back: before the next checkpoint, or overnight?
- Which evals are cheap (log-likelihood multiple choice) and which are expensive (long reasoning, pass@k, agentic)?
- Who adds benchmarks: one eval team, or every researcher? (This affects the plugin API and versioning.)
Scale estimates
- Assume 150 "cheap" benchmarks × 2,000 items × ~500 tokens ≈ 150M tokens, mostly prefill or short generation.
- Plus 50 reasoning or code benchmarks × 500 items × 8 samples × 8K tokens ≈ 1.6B generated tokens. These dominate.
- For a 70B checkpoint at ~1K generated tokens/s per GPU: 1.6M GPU-seconds ≈ ~450 GPU-hours per checkpoint. At 10 checkpoints a day across runs, that is 4,500 GPU-hours a day, about 190 GPUs running continuously. Tiering is essential.
Component walkthrough
- Trigger and tiering. Subscribe to checkpoint-commit events. A per-run policy decides the suite: a fast tier (minutes, log-likelihood and short generation) on every checkpoint, a full tier on every Nth checkpoint and on candidates for promotion. The Acme study found that decoupling evaluation scheduling from training cut evaluation makespan by 1.3–1.8× in its clusters Hu+ 2024.
- Work decomposition. Spin up an inference server per (checkpoint, hardware slice) and fan out item-level or chunk-level tasks through a queue. Item-level results make retries cheap and let several benchmarks share one warm server.
- Harness and reproducibility. Every result is keyed by (checkpoint, benchmark version, prompt template, few-shot set, sampling params, harness version, judge model version). Small formatting differences change scores measurably. OLMES argues for standardizing these choices Gu+ 2024. Open harnesses to build on or borrow from: lm-evaluation-harness and Inspect, whose datasets, solvers and scorers plus pluggable sandboxes map well to a platform plugin API.
- Sandboxed code eval. Reuse the RL environment service: no network, resource caps, per-language images, timeouts. Treat a timeout or crash as an explicit outcome category, not a silent fail.
- LLM-judge workers. Pin the judge model version and prompt. A judge upgrade is a new eval version, and scores before and after aren't comparable without re-scoring. Re-scoring is cheap if raw outputs are stored. Track judge agreement with human labels on a calibration set (see B5 and A8).
- Result store and statistics. Store item-level outputs and scores. That allows confidence intervals (questions as samples from a population Miller 2024), paired comparisons between checkpoints on the same items (far tighter than comparing two independent means), and item-level diffing to see what broke.
- Contamination checks. Keep an n-gram index of all benchmark items and check new training data against it at data-pipeline time (Case 1). At eval time, flag benchmarks with high overlap. Llama 3 reports an 8-gram overlap analysis and notes that for some benchmarks the method flags too much to be useful Llama Team 2024.
- Regression alerts. Alert when a score drops by more than k standard errors relative to the run's trailing trend or a baseline run, and stays down for two consecutive checkpoints, to filter noise.
Failure modes
- Eval infra bugs mistaken for model regressions: a tokenizer or chat-template change, a sandbox image update that breaks a dependency. Canary every harness change against frozen reference checkpoints whose scores are known.
Mention cost tiering and statistics together. "We can't afford the full suite on every checkpoint, and even if we could, a 1-point change on a 500-item benchmark is inside the noise. So: fast tier with paired comparisons on every checkpoint, full tier on milestones, alerts that are CI-aware." That one sentence covers both the infra and the science concerns.
- A8 · Evaluation: contamination, pass@k, judges, error bars.
- Inspect and lm-evaluation-harness: harness architectures worth borrowing.
- OLMES: standardizing evaluation choices for reproducibility.
- Adding Error Bars to Evals: the statistics behind CI-aware alerts.
Case 8 · Human preference / RLHF data collection platform
Prompt: "Design the platform that collects pairwise human preference data for reward-model training: task generation, sampling model responses, annotator workflow, quality control, and delivery to training."
Clarifying questions
- Data type: pairwise preference, graded ratings (Llama 2 used a four-level margin scale: significantly better, better, slightly better, negligibly better or unsure Touvron+ 2023), rankings of K responses (InstructGPT ranked 4–9 Ouyang+ 2022), or rubric-based multi-attribute scores (HelpSteer2 Wang+ 2024)?
- Annotator pool: an in-house expert team, a vendor workforce, or both? Domain experts (code, medicine) cost several times more than generalists.
- Volume and cadence. Llama 2 collected about 1.4M Meta preference comparisons in 14 weekly batches Touvron+ 2023. Anthropic's HH work iterated on a weekly cadence with fresh models Bai+ 2022. Is the loop online (new model each round) or one-shot?
- Safety: will annotators see harmful content? That brings wellness protections, opt-outs and content warnings.
Scale estimates
- Target 500k comparisons per quarter. At ~3 minutes per comparison that is 25k annotator-hours. At ~520 productive hours per annotator per quarter, about 50 full-time annotators. At $30–80 per hour loaded cost: $0.75–2M per quarter. Human labour, not GPUs, is the dominant cost.
- QC overhead: about 5–10% gold items and 10–20% multi-annotated items for agreement, adding 15–30% to labour.
Component walkthrough
- Task generation. Stratify prompts by domain, difficulty and language against a coverage target. Pick response pairs deliberately: current policy vs previous checkpoint, two samples from the same policy, policy vs strong reference. The pairing strategy decides what the reward model learns to distinguish. Randomize left/right order and record it, to measure position bias.
- Queueing and assignment. This is a work queue with leases (annotators abandon tasks), skill-based routing (code tasks to code-qualified annotators), and quotas so no single annotator dominates a slice. Prioritize tasks where the current reward model is uncertain (active sampling), which gives more information per labelled pair.
- Annotator UI. Side-by-side responses, a margin scale, optional rubric dimensions (helpfulness, honesty, harmlessness), a free-text rationale, flags (unsafe, both bad, prompt invalid), and time tracking.
- Quality control. (1) Gold items with known answers mixed in unannounced, measuring per-annotator accuracy. (2) Overlap items labelled by several annotators to compute inter-annotator agreement, with Cohen's kappa for pairs Wikipedia or Krippendorff's alpha for many raters and missing data Wikipedia. (3) Behavioural signals: implausibly fast completions, always picking the left response, copy-pasted rationales. (4) Review tiers: senior reviewers audit a sample and adjudicate disagreements. Low-agreement slices often mean the guidelines are ambiguous, not that annotators are bad.
- Storage and versioning. Append-only judgment events are the source of truth. Resolved labels come from an aggregation policy (majority vote, annotator-reliability weighting, gold-based filtering) that is itself versioned. A dataset snapshot records the policy version and the event range. Link every reward-model checkpoint to the snapshot it was trained on.
- Feeding reward-model training. Each batch triggers RM training and evaluation on a held-out human-agreement set. Track RM accuracy by slice.
Failure modes and trade-offs
- Annotators using LLMs to answer: detect it through timing, rationale style, and gold items designed to catch model-typical errors. Address it in contracts and tooling (paste restrictions).
- Trade-off: more annotators per item gives cleaner labels but fewer items. Reward-model accuracy usually benefits more from breadth plus targeted overlap than from labelling everything three times. HelpSteer2 shows that about 10k high-quality pairs can train a strong reward model Wang+ 2024, so quality can beat volume.
The design detail that impresses: "store judgments, not labels". The label is a derived view, computed by a versioned aggregation policy over immutable events. That lets you drop a bad annotator's contributions retroactively, change the aggregation rule, and reproduce exactly which data trained reward model v12.
- Llama 2 §3.2: a detailed public account of preference collection at scale (weekly batches, margin labels, safety vs helpfulness).
- InstructGPT: labeler selection, rankings of K responses, agreement analysis.
- Anthropic HH-RLHF: iterated online collection.
- HelpSteer2: small, high-quality, multi-attribute preference data.
- A4 · Post-training: how the data is consumed by RM and DPO training.
Case 9 · Synthetic data generation pipeline
Prompt: "Build a pipeline that generates 100B tokens of high-quality synthetic training data (reasoning traces, code with tests, instruction data), filtered and decontaminated, with full provenance."
Clarifying questions
- What kinds of data, and how is each verified? Code and math can be checked by execution or answer matching. Open-ended text needs judges or reward models.
- Which generator models, and what licences and terms of use apply to their outputs?
- Diversity targets: topics, difficulty, formats, languages. How do you seed diversity?
- Will the data be used for pre-training (volume) or post-training (quality)?
Scale estimates
- Target 100B kept tokens. With a ~25% pass rate after verification, judging and dedup, you must generate ~400B tokens.
- With a 70B generator at ~1K tokens/s per GPU (long-form decode at large batch): 4×10¹¹ ÷ 10³ = 4×10⁸ GPU-s ≈ 110k GPU-hours ≈ $275k. With an 8B generator at ~3–4K tokens/s per GPU: about 30k GPU-hours. Generator choice is mostly a quality-per-dollar question. Some pipelines use a strong model for hard seeds and a cheaper one for volume.
- Reference points: Cosmopedia generated over 30M files and 25B tokens with Mixtral-8x7B-Instruct Hugging Face 2024. Nemotron-4 340B reports that over 98% of the data used in its alignment process was synthetic NVIDIA 2024.
Component walkthrough
- Seeding for diversity. Synthetic data collapses toward the generator's favourite modes unless seeds force variety: topic taxonomies, real web documents as grounding, personas, difficulty levels. Self-Instruct, for example, kept a new instruction only if its ROUGE-L similarity to every existing one was below 0.7 Wang+ 2022.
- Generation. Reuse the batch inference platform (Case 6): sharded work units, leases, spot capacity. Generate K candidates per seed and pick the best (rejection sampling). Llama 3's post-training sampled typically 10–30 outputs per prompt and selected with a reward model, and used PagedAttention with shared prompt KV for over 2× throughput during rejection sampling Llama Team 2024.
- Verification. Execution-based where possible: run generated code against generated tests, but beware tests that agree with buggy code. Use cross-checks, such as solutions from several samples agreeing, or tests validated against a reference. Check final answers for math.
- Judging. Rubric-based LLM judges or reward models score quality. Store scores and pick thresholds later. Calibrate judges against human labels on a sample.
- Dedup and diversity control. Near-duplicate removal (MinHash as in Case 1), plus embedding clustering with per-cluster caps so one template doesn't dominate.
- Decontamination. Check against every eval set in the benchmark registry with n-gram and embedding similarity. Synthetic generators can regurgitate benchmark items they memorized, which makes this more important here than for web data.
- Provenance. Record per-sample metadata (generator, version, template, verifier, judge, licence). If a generator version turns out to be flawed, or a licence term changes, you can filter its outputs out of future dataset versions.
Failure modes and trade-offs
- Model collapse and homogenization: training repeatedly on model outputs narrows the distribution. Shumailov+ showed collapse when training recursively on generated data Shumailov+ 2024. Mix with real data, track diversity metrics (distinct n-grams, embedding spread), and cap the synthetic share per domain.
- Judge exploitation: generators produce what judges like (length, confident tone). Use several judges and length-controlled scoring, and spot-check with humans.
- Trade-off: a bigger generator gives higher pass rates but more GPU-hours per token. Compute cost per kept token, not per generated token.
Frame synthetic data as "batch inference plus a filtering funnel". The cost metric is GPU-hours per accepted token, which depends on pass rate. Improving the verifier or the seed quality often beats buying a bigger generator.
- Cosmopedia: a hands-on account of seeding and generating 25B synthetic tokens.
- Nemotron-4 340B report: a synthetic-data-heavy alignment pipeline.
- Textbooks Are All You Need (phi-1): synthetic "textbook" data for code.
- AI models collapse when trained on recursively generated data: the risk to design against.
Case 10 · Feature store and training pipeline for a recommender with LLM-generated features
Prompt: "Our recommender ranks items for 100M users. We want to add LLM-generated item features (topic tags, quality scores, embeddings) and keep training and serving consistent." This classical-ML design still comes up, now with an LLM twist.
Clarifying questions and estimates
- Catalogue size and churn: assume 10M items and 200k new items a day. Request rate: assume 20k ranking requests/s at peak, each scoring 500 candidates.
- LLM item features: 10M items × ~1,500 input tokens with an 8B model → 1.5×10¹⁰ prefill tokens ≈ 2.4×10²⁰ FLOPs ≈ ~250–350 GPU-hours for the full backfill, a few thousand dollars. The daily 200k new items are trivial. The cost is in the backfill whenever the prompt or model changes.
- Online lookups: 20k requests × 500 candidates = 10M item-feature reads/s. That is too many for a remote key-value store per item, so item features must be cached in the ranking service's memory or co-located. 10M items × ~2 KB ≈ 20 GB fits in RAM per replica.
Key concepts
- Online/offline consistency. Define each feature once (code and version) and materialize it to both stores from the same pipeline. Uber's Michelangelo post describes this split: an offline store on Hive/HDFS for training and an online store on Cassandra for serving Uber 2017. Its feature store is now known as Palette Uber 2024. Open-source equivalent: Feast.
- Point-in-time correctness. Each training example must use feature values as they were at the label's event time, never later. Otherwise you leak the future (for example, an item's click count including the click you're predicting). Feast documents point-in-time joins with TTLs relative to each row's timestamp Feast docs. Table formats with time travel (Iceberg's
VERSION AS OF/TIMESTAMP AS OF) help with reproducibility Iceberg docs. - LLM features add versioning pain. Changing the prompt or model changes every item's features. Treat
llm_tags_v3as a new feature column, backfill it fully (cheap, per the estimate above), train a new model on it, and A/B test. Never overwrite v2 in place while a model trained on v2 is serving. - Log-and-wait. The most robust way to avoid training/serving skew is to log the exact feature vector used at serving time and join labels to it later. Then training sees exactly what serving saw.
Computing an LLM feature from item text that is updated later (title edits, new reviews) and joining the current LLM output onto historical training examples. That leaks future information. Version the LLM features with their computation time and join point-in-time like any other feature.
- Meet Michelangelo (Uber, 2017): the canonical end-to-end ML platform write-up.
- Feast: point-in-time joins.
- Metaflow (Netflix): a workflow framework for ML pipelines with versioned artifacts.
Part 3 · Cross-cutting concerns
Reproducibility
- Everything that defines a run is versioned: code commit plus container image digest, the fully resolved config (no "latest" anywhere), dataset snapshot IDs, the tokenizer version, and the data mixture spec.
- Seeds and data order: a deterministic sampler (§1.2) means the batch at step \(s\) is a pure function of the config. Save RNG states (Python, NumPy, torch, CUDA) in checkpoints.
- Bitwise vs statistical reproducibility: bitwise reproducibility across different GPU counts or kernel versions is usually unrealistic (non-deterministic reductions, autotuned kernels). Aim for "same config gives a statistically indistinguishable loss curve", and reserve bitwise determinism for debugging small runs.
- Evals are part of reproducibility: pin the harness, prompts and judge versions (Case 7).
Lineage and governance
The lineage graph connects source data → processed dataset versions → training runs → checkpoints → eval reports → registered models → deployments. Use it to answer: "which deployed models were trained on data from source X?" (takedowns, licence changes, contamination findings), "what changed between model v11 and v12?", and "who approved this release?". Emit lineage events from every job (OpenLineage-style) and keep the graph in a queryable store. Governance layers on top: dataset licences and usage restrictions as metadata, access control on sensitive datasets, and release approvals recorded in the registry.
Cost optimization
| Lever | Where it applies | Caveat |
|---|---|---|
| Spot / preemptible | Batch inference, evals, data processing, small experiments | Needs small idempotent work units. Rarely used for large synchronous training |
| Raise utilization | Shared clusters (Case 3): scavenger tier, idle reaping, backfill | Preemption costs must be tracked |
| Right-size hardware | Small models and encoders on cheaper GPUs; decode on memory-bandwidth-rich parts | Benchmark tokens per dollar on your workload |
| Raise MFU | Training: parallelism tuning, FP8, kernel fusion (A3, A7) | Engineering time vs GPU savings |
| Raise goodput | Large runs: faster restart, async checkpoints (Case 2) | Most valuable at the largest scale |
| Cache and reuse | Prefix caching in serving, stage caching in data pipelines, eval result caching | Cache keys must include every version that matters |
| Reduce work | Smaller proxy models for data ablations, tiered evals, early stopping of bad runs | Proxy results must correlate with the target scale |
Security of model weights
Frontier weights are among an AI company's most valuable assets, and a target for well-resourced attackers. RAND's report on securing model weights identifies about 38 distinct attack vectors and proposes five security levels (SL1–SL5), matched to attacker capability Nevo+ 2024. Design-level controls to mention:
- Least-privilege access to weight storage. Few humans have direct read access. Access goes through audited services.
- Encryption at rest and in transit, with keys in KMS or HSMs. Training and serving nodes get short-lived credentials scoped to specific artifacts.
- Egress monitoring and limits on clusters holding weights (a multi-TB transfer should be impossible or loud).
- Supply-chain hygiene: signed artifacts, and safe serialization (safetensors instead of pickle, which can execute arbitrary code on load safetensors).
The training-to-serving handoff
- Evaluate the artifact you serve, not the checkpoint you trained. Quantization, format conversion and engine kernels can each change outputs. Chat-template mismatches are a classic silent regression.
- Eval gates are policy as code: thresholds per benchmark and per safety category, with explicit human override recorded in the registry.
- RAND: Securing AI Model Weights: attack vectors and security levels.
- OpenLineage: an open standard for lineage events.
- MLflow Model Registry: a reference design for versions, stages and approvals.
- PyTorch Distributed Checkpoint: resharding checkpoints between training and export layouts.
Interview question bank
A 4k-GPU job's goodput dropped from 90% to 60%. How do you debug it?
First, split goodput into its two parts. Is the job up less of the time, or is it slower while up? The dashboard's time breakdown (train, checkpoint, restart, hang, waiting for nodes) answers that in minutes.
- More restarts: look at the failure taxonomy. Common causes are one node that keeps failing and rejoining (quarantine it), failures correlated by rack or leaf switch, a new hardware batch, or a recent NCCL, driver or code change.
- Slower recovery: find where the restart minutes go. Usual suspects are an exhausted spare pool, image pulls, NCCL init, and checkpoint loads from degraded storage.
- Slower steps: check per-rank timings. One persistently slow rank points to hardware. Variance that correlates with the data or with GC points to software. Also check whether checkpoint duration or data-loader wait time has grown.
Fix the biggest bucket first, then confirm the improvement on the dashboard.
How often should you checkpoint? Derive it.
Define three quantities: blocking cost \(C\) per checkpoint, job MTBF \(M\), and interval \(\tau\).
- The overhead fraction is \(C/\tau + (\tau/2 + R)/M\), where \(R\) is the restart time.
- Setting the derivative to zero gives \(\tau^* = \sqrt{2CM}\) (Young; Daly refined it).
- Example: \(C\) = 30 s and \(M\) = 3 h give \(\tau^* \approx\) 13 min.
Two follow-ups. With async checkpoints, \(C\) shrinks to seconds, so \(\tau^*\) drops to minutes and the write bandwidth of durable storage becomes the limit. \(R\) doesn't depend on \(\tau\) and often dominates, so once checkpoints are async, the next thing to optimize is restart time.
Design dedup for 100T tokens.
That's about 10¹¹ documents. Work in tiers:
- Exact dedup: group by URL and by content hash.
- Fuzzy dedup: MinHash-LSH. Signatures are a map stage that reads ~400 TB once. Emitting 14 band keys per document means a shuffle of about 22 TB. Within each bucket, emit star edges to avoid quadratic pairs. Then run connected components and keep one document per cluster, chosen by policy.
- Boilerplate (optional): line-level dedup with Bloom filters or suffix arrays.
Engineering points to raise: skew from giant boilerplate buckets (cap or salt them), spilling to local SSD, and idempotent stages. Treat the dedup scope (per-snapshot vs global) and the similarity threshold as choices to settle with small-model ablations.
How do you keep 2,000 sandboxes warm for RL rollouts?
Size the warm pool with Little's law. 2,000 concurrent episodes of 300 s means about 7 new sandboxes per second. With a 30 s cold start, you need at least 200 warm sandboxes per image mix, plus headroom for bursts. A pool controller per image class keeps them ready and refills the pool as sandboxes are claimed.
Then make cold starts cheaper:
- Pre-pull images onto nodes.
- Build per-repo images rather than per-task images.
- Use a lighter runtime (crun with pre-created cgroups, as in Kimi k1.5).
- Snapshot and restore fully set-up environments.
Give every sandbox a TTL and a heartbeat, and run a reaper for orphans. When the pool runs dry, apply back-pressure to rollout workers. Letting them time out produces fake zero rewards.
How big is a checkpoint for a 70B model, and how long does it take to save?
The size depends on what you save:
- Weights only: 140 GB in bf16.
- Full training state: about 14–16 B per parameter (fp32 master weights plus two Adam moments), so roughly 1 TB. Spread across 512 GPUs, that's about 2 GB per GPU.
A synchronous save to storage with 200 GB/s aggregate bandwidth takes about 5 s in theory and tens of seconds in practice. With async checkpointing, training blocks only for the GPU→host copy, which takes well under a second. The restore path is often more critical, because the whole job sits idle while it loads. Keep recent checkpoints in host memory, peer memory or on local NVMe.
How would you detect and handle a hang in a 10k-GPU job?
Hangs are silent, so you have to detect them actively.
- Detect: every rank sends a step heartbeat. Declare a hang when global progress stalls for a few multiples of the normal step time, rather than waiting for the default collective timeout.
- Diagnose: dump Flight Recorder traces and stack traces from every rank. Find the ranks that never entered the pending collective, or whose sequence numbers diverge.
- Confirm: cross-check the suspects against DCGM, XID errors and NIC counters.
- Recover: quarantine the suspect nodes, swap in spares and restart from the last checkpoint.
Track hang frequency and time-to-detect as goodput metrics.
Size the hot spare pool for a 2,048-node training job.
Use Little's law: nodes out of service ≈ failure rate × repair time. About 8 node failures a day with a 2-day repair turnaround means roughly 16 nodes are out at any moment. Add headroom for correlated failures (a whole rack or switch), which gives about 24–40 spares, or 1–2% of the job.
Spares should be burned in, have images pre-pulled, and sit where swapping them in keeps the job's topology tight. Compare the cost of 2% idle spares with the cost of all 2,048 nodes idling while a restart waits for hardware.
Design the autoscaler for an LLM serving pool. What signal do you scale on?
Don't scale on CPU, and don't scale on raw GPU utilization either: a GPU reads busy whether it's comfortable or overloaded. Scale on:
- queue depth per priority class;
- KV-cache utilization;
- TTFT and ITL percentiles measured against the SLO.
Cold starts take 1–5 minutes, so set a target below saturation, scale up early, pre-warm on the daily traffic pattern, and scale down slowly with hysteresis. To shorten cold starts, cache weights on local NVMe, use streaming loaders, or copy weights GPU-to-GPU from a running replica. For long-tail models, use scale-to-zero only if their SLO tolerates a cold start. Otherwise consolidate them as LoRA adapters on a shared base model.
How would you serve 500 fine-tuned variants of the same 8B model cheaply?
If they are LoRA fine-tunes, serve them all on shared base-model replicas with multi-adapter batching (the Punica / S-LoRA idea, available in vLLM). Each adapter is tens of MB, so all 500 fit in host RAM and the working set fits in GPU memory. The platform needs three pieces:
- an adapter cache hierarchy (GPU → host RAM → object store);
- adapter-affinity routing, so requests land where their adapter is already loaded;
- admission control for requests that need an adapter swap.
For full fine-tunes, convert or distill them to LoRA, give dedicated replicas only to the hot ones, and put the rest on scale-to-zero.
In RL training, the trainer GPUs sit idle 50% of the time. What do you do?
Find out what the trainer is waiting for.
- Long-tail generation: a few long rollouts hold up each batch. Options are async generation with a staleness bound η plus importance correction, partial rollouts, more rollout capacity, or a different rollout/trainer GPU ratio.
- Slow generation overall: check engine batch occupancy, prefix caching across group samples, and environment latency.
- Slow weight sync: move to bucketed NCCL or RDMA broadcast with shard-aware resharding.
Measure idle percentage per pool, staleness and reward curves before and after the change. Async speedups can be eaten by off-policy effects.
How do you make weight sync between a 70B trainer and 128 inference nodes fast?
Each sync moves about 140 GB. Fanning that out through object storage would mean 18 TB of reads, so broadcast over the network instead. Reshard from the trainer layout to the engine's TP layout, then stream buckets through a pipelined tree or ring broadcast over InfiniBand or RoCE. That takes a few seconds end to end.
Further improvements:
- Overlap the transfer with the tail of generation.
- Swap weights atomically by version, or in flight.
- Use CUDA IPC when trainer and engine are colocated.
- Sync only the adapter when training LoRA.
Tag every engine with its weight version, so each trajectory records which policy produced it.
Design an eval platform that runs 200 benchmarks on every checkpoint without bankrupting us.
Tier it:
- Fast suite on every checkpoint: log-likelihood and short-generation tasks.
- Full suite only on milestones and release candidates: long reasoning, pass@k, agentic and judge-based evals. These cost hundreds of GPU-hours per 70B checkpoint.
Run evals on the scavenger tier with a guaranteed minimum share. Use one inference server per checkpoint and fan out item-level tasks to it. Cache results by (checkpoint, benchmark version, harness version, judge version). Store item-level outputs, so you can re-score cheaply, compute confidence intervals and compare checkpoints with paired tests. Alert only on CI-aware regressions that persist across two checkpoints.
How do you guarantee a training run is reproducible six months later?
Pin everything:
- image digest and code commit;
- the fully resolved config;
- immutable dataset snapshot IDs;
- tokenizer version and mixture spec;
- seeds.
The data loader should be deterministic, checkpoints should include RNG states, and the lineage store should link the run to its inputs and outputs. Retain the datasets, or keep enough to regenerate them. Be honest about the limits: bitwise reproducibility across different hardware or GPU counts is usually out of reach, so aim for statistical reproducibility, plus bitwise reproducibility for small debug runs.
Your preference-data reward model stopped improving even though you keep adding labels. What do you check?
Work through the likely causes in order:
- Label quality: inter-annotator agreement over time, gold-item accuracy, signs of LLM-assisted annotation, and drift in the guidelines.
- Informativeness: new pairs may be too easy, or drawn from models the RM already separates well. Switch to active sampling of pairs where the RM is uncertain, using current-policy outputs.
- Coverage: new prompts may be piling into categories that are already saturated.
- The eval set itself: the RM's ceiling is the human agreement rate.
Because judgments are stored as immutable events, you can re-aggregate them with stricter filters and retrain to test each hypothesis cheaply.
How do you prevent eval benchmarks from leaking into training data?
Make decontamination a pipeline invariant:
- Keep a benchmark registry with n-gram and embedding indexes over every eval item.
- Run a decontamination stage in every data pipeline, web and synthetic, and record its version in the lineage graph. Synthetic generators can regurgitate memorized items, so they need it too.
- Report overlap statistics per benchmark at eval time, as in Llama 3's 8-gram analysis, and base decisions on fresh or private held-out sets.
- When you add a new benchmark, re-scan past training datasets so you know which runs saw it.
Where would you use Kafka in an ML platform, and where would you not?
Use it for high-volume, append-only streams that have many consumers or need replay:
- inference request logs, which feed monitoring, eval sampling, billing and (with consent) training data;
- interaction events for recommender features;
- annotation events;
- possibly trajectories in a disaggregated RL system.
Don't use it on the training hot path: gradients move through NCCL. Don't use it as a batch work queue either. Per-item leases and retries fit a database-backed queue or an SQS-style service better than partition-ordered consumption.
A model looks great in offline evals but users report worse answers after deploy. What could have gone wrong between training and serving?
Common causes:
- a chat-template or tokenizer mismatch;
- quantization degrading capabilities the eval gate didn't cover;
- different engine numerics;
- different default sampling parameters;
- system-prompt differences;
- context truncation in serving.
The fix is process. Evaluate the exact serving artifact and configuration, run shadow traffic with judge-scored comparisons before the canary, and record the full serving config in the registry next to the model version.
Estimate the GPU-hours to pre-train a 70B dense model on 15T tokens, and how many GPUs you'd want.
The compute estimate:
- FLOPs ≈ 6 × 70×10⁹ × 15×10¹² = 6.3×10²⁴.
- At ~40% MFU on H100, about 400 TFLOPS effective per GPU, that's about 4.4M GPU-hours, or $9–18M.
- Divide by about 0.9 to account for goodput losses.
- To finish in about 60 days you need 4.9M ÷ 1,440 h ≈ 3,400 GPUs, so plan on 4k GPUs including spares.
Then ask whether the run is meant to be compute-optimal or deliberately over-trained for inference efficiency (see A2).