Track C · Design interviews

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

MetricDefinitionWhat 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
GoodputUseful training progress (tokens or steps that survive into the final model) ÷ what ideal, uninterrupted hardware would deliver in the same wall timeFailure rate, detection time, restart time, checkpoint overhead, lost work since the last checkpoint, stragglers
Effective training timeFraction of wall-clock time spent training, not restarting or recovering. Llama 3 reported above 90% Llama Team 2024Same as goodput, minus the MFU term
Allocation vs utilizationAllocation: GPUs assigned to jobs ÷ GPUs owned. Utilization: SM activity or MFU of those allocated GPUsScheduler policy, fragmentation, idle interactive sessions, data-loader stalls
GPU-hours and $The budget unit for every jobAll of the above, plus hardware choice, spot or preemptible capacity, right-sizing
Intuition

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

BottleneckSymptomHow you see itTypical fixes
ComputeHigh SM activity, MFU near the hardware's realistic ceilingProfiler shows GEMM-dominated timelineLower precision (FP8), better kernels, more GPUs
HBM capacityOOM, forced small micro-batches, heavy activation recomputeMemory snapshots, allocator statsSharding (FSDP/ZeRO), activation checkpointing, offload, more TP/PP (A3)
HBM bandwidthLow arithmetic intensity, e.g. decodeRoofline position, DRAM throughput countersBigger batches, quantization, KV-cache compression (A6)
InterconnectGPUs idle inside collectives. Step time grows with scaleNCCL timings, exposed communication in traces, link countersOverlap comm with compute, topology-aware placement, change the parallelism layout
Storage I/OSlow checkpoint save or load, slow startupFS or object-store throughput, time-to-first-stepAsync or sharded checkpoints, local caches, parallel reads
Data loaderGPU gaps between steps, CPU pegged on preprocessingData-wait time per step, host CPU usagePre-tokenize, more workers, prefetch, move preprocessing offline
Scheduler and queueJobs wait while GPUs sit idle, or big jobs never startQueue time, allocation rate, fragmentationGang plus backfill, topology-aware packing, preemption, quotas with borrowing
StragglersStep time set by the slowest rank. p99 rank time ≫ p50Per-rank step timing, NCCL wait attributionDetect and evict slow nodes, fix thermals or links

An answering framework

1. Clarify What is the workload (train / serve / batch / data)? Scale (GPUs, tokens, QPS, docs)? Objective: throughput, latency, cost, or time-to-result? Hard constraints (budget, region, compliance)? 2. Estimate FLOPs (6ND for training, 2N per token for inference), bytes (data, checkpoints, KV), GPU-hours, $, failure rate (MTBF ≈ per-GPU MTBF / N). Say every assumption out loud. 3. Skeleton Control plane (API, scheduler, metadata DB, state machine) vs data plane (GPUs, storage, network). Draw the data flow first, then the control flow. 4. Bottleneck Which resource saturates first? Use the taxonomy. Size everything else around it. 5. Failure What fails, how often, how is it detected, what is the blast radius, how do you resume? Idempotency + checkpoints + leases + retries. 6. Ops Observability, alerting, lineage and reproducibility, multi-tenancy, security, cost. 7. Trade-offs What you'd change at 10x scale or 10x lower budget. What you deliberately left out.
Interview angle

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.

QuantityBallparkUse it for
H100 SXM BF16 dense~990 TFLOPS (FP8 ~1,980)Training time from 6ND. Realistic MFU 35–50%
H100 HBM80 GB, ~3.35 TB/sDecode speed, KV capacity
B200-class GPU~2.2 PFLOPS BF16 dense, ~180 GB HBM3e, ~8 TB/sNext-gen capacity planning
NVLink (per GPU, total)H100 ~900 GB/s. Blackwell ~1.8 TB/sTP inside a node or NVL72 domain
InfiniBand / RoCE per NICNDR 400 Gb/s ≈ 50 GB/s. 8 NICs per node ≈ 400 GB/s/nodeDP, PP, weight broadcast across nodes
PCIe Gen5 x16~64 GB/sHost↔GPU copies (async checkpoint staging)
Local NVMe~5–12 GB/s per driveShard caches, staged checkpoints, model weights
Parallel FS / storage cluster aggregateHundreds of GB/s to several TB/s. Llama 3's storage: ~2 TB/s sustained, ~7 TB/s peakCheckpoint save/restore time
Object storePer-prefix request limits (S3: 3,500 PUT / 5,500 GET per second per prefix). Throughput scales with parallel connections. ~100 MB/s per streamDatasets, durable checkpoints, model registry
Training state bytes per paramBF16 weights 2 B. Mixed-precision Adam training state ~14–16 BCheckpoint 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 tokenBatch-job and prefill cost
Job MTBF at scaleLlama 3: ~3 h at 16k GPUs. Scales ~1/NCheckpoint interval, spare pool size
Bytes per token (text)~4 bytes raw text. 4 bytes stored as uint32 token IDsCorpus sizing
GPU price~$2–4 per H100-hour. Spot or preemptible often 30–70% cheaperCost estimates
Spot / preemptible noticeAWS 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
May be out of date

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.

Go deeper

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:

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

ConceptWhat it meansDesign note
Gang scheduling (all-or-nothing admission)Admit a job only when all its pods or nodes can start togetherWithout 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 placementPack a job into the fewest switches, NVLink domains or blocksEncode topology as node labels (rack, leaf, block). Prefer tight packing for large jobs and leave whole blocks free
Common mistake

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 (S3 / GCS / Azure Blob / on-prem: Ceph, Tectonic-like) Durable source of truth · PB–EB · cheap per GB · high aggregate throughput, high per-request latency · datasets, checkpoints, registry Shared fast tier: parallel FS (Lustre, WEKA, GPFS, 3FS) or a distributed cache POSIX semantics · 100s of GB/s to TB/s · hot datasets, checkpoints in flight, shared code and envs · expensive per GB Node-local NVMe + host RAM + GPU HBM Fastest, ephemeral, lost with the node · shard caches, staged checkpoints, model weights for serving, page cache read / prefetch async persist
Three tiers. Data flows down for reading. Checkpoints flow up asynchronously so the GPUs don't wait on the slow tier.

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:

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.

Interview angle

"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 SparkRay DataDask
ModelBulk-synchronous stages over partitioned DataFrames or RDDs. Shuffle is a first-class operationStreaming execution over blocks, with heterogeneous CPU and GPU operators in one pipeline Ray docsTask graphs over pandas/NumPy-like collections in Python
Best atHuge joins, group-bys, dedup shuffles, SQL over PB-scale tables, mature fault tolerancePipelines mixing CPU preprocessing with GPU model inference (classifiers, embedding, batch LLM), and streaming into trainingPython-native analytics at medium scale, research workflows
Weak atGPU stages (possible but awkward), Python UDF overheadFor the very largest joins and shuffles, Spark is the more widely deployed path. Benchmark Ray on your own data before committingVery large shuffles, operational maturity at PB scale
Typical useDedup, joins with URL blocklists, aggregationQuality-classifier scoring, batch inference, tokenization feeding trainingAd-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?

NeedTypical toolsDesign core
Experiment trackingWeights & Biases, MLflow, in-houseRun ID → config, code commit, metrics time series, artifacts. Metrics go to a time-series store. Run metadata goes to a relational DB
Model registryMLflow Model Registry, in-houseModel name → versions → stage (candidate, staging, prod), with links to training run, eval report and approvals
Dataset versioningApache Iceberg / Delta Lake snapshots, lakeFS, DVCImmutable snapshots with IDs. A training job pins a snapshot ID, never "latest"
LineageOpenLineage, in-house graphA 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:

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

Intuition

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.

Go deeper

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

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:

Common Crawl WARC, ~100 snapshots URL filter blocklists, opt-outs Extract text trafilatura / resiliparse Language ID fastText lid Heuristic filters length, repetition, lines Dedup (shuffle-heavy) exact: URL + content hash fuzzy: MinHash → LSH bands → group-by → connected components Annotate (GPU + CPU) quality score, toxicity, PII spans, topic/domain labels Select + transform thresholds, PII redaction, decontamination vs eval sets Tokenize + shard uint32 IDs, ~100s MB shards Mixture spec weights per source/bucket Deterministic loader f(seed, step, rank) → sample Object store: Parquet per stage table snapshots (Iceberg/Delta) keyed by (stage, stage_version, input_snapshot) Lineage + manifest DB doc_id → source URL, crawl, stage versions, shard → doc ranges; takedown index every stage writes outputs + lineage records
Stages are idempotent batch jobs over partitioned tables. Annotation stages add columns and never delete rows, so changing a threshold re-runs only selection and the stages after it.

Component walkthrough

  1. 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.
  2. 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.
  3. 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.
  4. Heuristic filters. Record each rule's result as a boolean or score column. Don't drop rows yet.
  5. Dedup. Exact (normalized-content hash), fuzzy (MinHash-LSH), and optionally line-level or substring-level. See the deep dive below.
  6. 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.
  7. 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).
  8. 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:

  1. 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).
  2. 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.
  3. 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.
  4. 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.

Intuition

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

Failure modes

Interview angle

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.

Go deeper

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

Scale estimates

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}$$
ScenarioC (blocking)Rτ* = √(2CM), M = 180 minOverhead≈ Availability
Naive: sync checkpoint, slow restart2 min25 min27 min7.4% + 21.4% ≈ 29%~71%
Sync checkpoint to fast storage, decent restart0.5 min12 min13 min3.8% + 10.3% ≈ 14%~86%
Async + in-memory checkpoint, hot spares0.05 min4 min4.2 min1.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.

Job controller desired state: run R, N nodes state machine: RUNNING → FAILED → DRAINING → RESTARTING → RUNNING (metadata DB: Postgres / etcd) Node health service DCGM, XID/ECC, NIC/link checks pre-flight burn-in, NCCL tests node states: healthy / suspect / quarantined / in repair Scheduler + spare pool N active nodes + S hot spares (burned-in, images pre-pulled) topology-aware swap-in Training ranks (16,384 GPUs · 2,048 nodes) Per-rank agent step heartbeat, per-rank timings Collective tracer flight recorder ring buffer Async ckpt HBM → host RAM (+ peer copy) Deterministic data loader f(seed, step, rank) heartbeats, failures Checkpoint tiers host RAM / peer RAM: every few minutes, seconds to restore local NVMe: staged, then uploaded parallel FS / object store: durable, every 30–60 min sharded (DCP-style), reshardable on load Run dashboard + alerting loss, grad norm, tokens/s, MFU, step-time p50/p99/max goodput breakdown, restarts, time-to-recover straggler heatmap per rank, node-health events pages: hang > 5 min, loss spike, MFU drop
A small consistent control plane (controller, health service, scheduler) drives a large failure-prone data plane. Checkpoints go through tiers so the GPUs never wait for durable storage.

Component walkthrough

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:

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:

May be out of date

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

Run dashboard and alerting

PanelSignalAlert
ProgressTokens/s, step time, MFU, tokens trained vs planMFU drop > 10% sustained for 15 min
Health of the modelLoss (train and held-out), grad norm, update/weight ratio, NaN countLoss spike > k σ, NaN, grad-norm explosion
Goodput breakdownTime in train, checkpoint, restart, hang, idle-waiting-for-nodesDaily goodput below target
StragglersPer-rank compute time heatmap, p99/p50 ratioSame node above threshold for N windows
FleetHealthy, suspect, quarantined and spare node counts. XID and ECC ratesSpare pool below minimum
CheckpointsLast committed step, save and upload durationsNo committed checkpoint in 2× expected interval
Interview angle

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.

Go deeper

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

Scale estimates

CLI / SDK / UI job spec + team + tier Admission + quota Kueue ClusterQueues / Slurm QOS Policy config (GitOps) team quotas, cohorts, priorities, limits Tier: production / reserved big pre-training runs guaranteed blocks, non-preemptible topology-tight placement Tier: team batch quota per team, can borrow borrowed share is preemptible fair share within team Tier: interactive + scavenger dev pool: ≤8 GPUs, idle reaper scavenger: lowest priority, fills gaps, evicted first Placement: gang admission · topology-aware packing · backfill · checkpoint-aware preemption node-granular for multi-node jobs, GPU-granular bin-packing for small jobs on a dedicated partition 512 nodes · 4,096 GPUs labelled by block / leaf / rack / health Accounting + utilization warehouse allocation, SM activity, MFU per job → chargeback
Tiers separate guarantees (reserved), elastic team quotas (borrowable, preemptible when borrowed) and best-effort work (interactive with reaping, scavenger). Accounting closes the loop.

Component walkthrough

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:

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

Interview angle

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.

Go deeper

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

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.

Task service prompt pools, curriculum, difficulty filtering Rollout fleet (inference engines) vLLM / SGLang replicas, policy version v continuous batching, prefix cache (G samples share prompt) partial rollouts: token budget per step agent loop driver per episode Environment service sandbox pool (containers / microVMs), warm pool per image exec API: run, read, write, reset CPU nodes, autoscaled tool call obs Reward service unit tests in sandboxes, answer checkers, RM / judge replicas Trajectory buffer tokens, masks, rollout logprobs, policy version, reward; staleness filter ≤ η Trainer (FSDP / Megatron) policy (+ reference logprobs, + critic) importance correction for stale data checkpoints every K steps Weight sync trainer layout → engine layout (reshard) NCCL broadcast tree / RDMA P2P / CUDA IPC version counter v+1 in-flight or between-batch swap Observability reward curves per env, length dist, staleness, gen/train idle %, sandbox errors, hacking alarms Async loop: generation never waits for training; trainer consumes data up to η versions stale.
An inference fleet, an environment service and a reward service feed a trajectory buffer. The trainer consumes it and pushes new weights back. This expands the A5 diagram with the platform services.

Component walkthrough

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:

Failure modes

Interview angle

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.

Go deeper

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

Scale estimates

Clients products, internal, batch API gateway auth, quotas, rate limits per tenant, SLO tier → priority, model alias resolve Model registry + deploy controller alias → version → artifact, engine config, canary %, autoscale policy (GitOps) Per-model router queue per priority, admission control prefix/KV-aware + load-aware pick LoRA-affinity, canary split Hot model pool prefill workers ⇄ decode workers (disaggregated, KV transfer) many replicas, TP within node Shared LoRA pool one base model, many adapters GPU / host-RAM / object-store adapter cache tiers Long-tail pool scale-to-zero / min=1, weights cached on local NVMe, fast loaders (streaming) Autoscaler signals: queue depth, KV-cache usage, TTFT/ITL vs SLO, tokens/s + schedule-based pre-warming Telemetry: per-request TTFT, ITL, tokens, cache hit, version, adapter → logs (Kafka) → metrics, eval sampling, billing
The gateway handles tenancy, the router handles placement, and pools are shaped by traffic: dedicated disaggregated pools for hot models, shared base-model pools for LoRA variants, scale-to-zero for the long tail.

Component walkthrough

Failure modes

May be out of date

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.

Go deeper

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

Scale and cost estimates

Corpus snapshot Parquet, snapshot ID S (Iceberg/Delta) Planner split into ~50k shards, ~2–5 min of work each Work-queue table (Postgres) shard_id, state, lease_owner, lease_expiry, attempts, output_uri GPU workers on spot / preemptible nodes (autoscaled) CPU: read + tokenize length-bucketed batches GPU: encode TEI / vLLM / custom Write shard output tmp path → atomic rename lease / heartbeat / complete Spot handler on interruption notice: stop taking work, flush or release lease (shard requeues) Object store: embeddings Parquet one file per shard + manifest (model_version, snapshot S, shard list, row counts, checksums) Offline index build → new collection bulk import / build IVF-PQ or graph index catch-up: CDC since S, embed deltas validate recall → flip alias (blue/green)
A lease-based work queue over a pinned snapshot, idempotent per-shard outputs, a manifest commit, then an offline index build and an alias flip. Spot interruptions only ever cost one shard's worth of work.

Component walkthrough

Failure modes and trade-offs

Interview angle

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.

Go deeper

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

Scale estimates

Checkpoint event manifest committed (run, step, URI) Eval orchestrator suite = f(run policy, step): fast tier every ckpt, full tier on milestones; dedup by cache key Benchmark registry versioned datasets + prompts + scorer code; access control; n-gram index for decontamination Generation workers ephemeral vLLM/SGLang servers per checkpoint, scavenger tier, item-level tasks via work queue Scoring workers exact match / regex / math checkers; code → sandbox pool; LLM-judge pool (pinned judge) Sandbox pool no network, CPU/mem/time caps, per-language images, shared with RL env service Result store item-level rows (prompt hash, output, score, sampling params) in columnar store; aggregates + CIs in SQL DB; key = (ckpt, bench_version, harness_ver) Dashboards + regression alerts score vs step per run, run vs baseline, CI-aware alerts (drop > 2 SE sustained), item-level diff viewer, contamination flags, promotion gates for model registry
Checkpoint events trigger tiered suites. Generation runs on ephemeral servers in the scavenger tier. Scoring fans out to checkers, sandboxes and judges. Item-level results support CI-aware alerts and diffs.

Component walkthrough

Failure modes

Interview angle

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.

Go deeper

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

Scale estimates

Prompt sources annotator-written, sampled traffic (consented, scrubbed), synthetic, red-team Task generator stratify by domain/difficulty, pick model pair (policy vs ckpt, temps), inject gold + overlap items Response sampler calls serving platform (Case 5) for N model variants; logs model version + sampling params Task queue + assignment skill-based routing, leases with expiry, overlap targets, priority by data need (active sampling) Annotator UI side-by-side (randomized order), margin scale + rubric + rationale, flags (unsafe, both bad), timing QC service gold accuracy, IAA (κ / α), time anomalies, reviewer escalation Label store (append-only events) judgment events: annotator, task, choice, margin, rubric, time, UI version → resolved labels (aggregation policy v) → versioned dataset snapshots Delivery to training dataset vN (snapshot ID) → reward model training → RM eval (held-out agreement) → new policy → next round's sampler dashboards: volume, coverage, quality
Store raw judgment events immutably. Resolved labels and dataset snapshots are derived by a versioned aggregation policy, so QC decisions can be revisited without re-annotating.

Component walkthrough

Failure modes and trade-offs

Interview angle

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.

Go deeper
  • 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

Scale estimates

Seeds topics taxonomy, web docs, personas, problem templates Generate batch engine (Case 6), K samples per seed, prompt templates vN Verify execute tests, check answers, format/schema validation (sandbox pool) Judge / score LLM judge or RM rubric scores stored as columns (annotate, don't delete) Dedup + diversity MinHash / embedding clustering, per-cluster caps Decontaminate n-gram + embedding match vs eval registry (Case 7) Select + mix thresholds, quotas per category Dataset vN snapshot + datasheet Provenance record per sample seed ID · generator model + version · prompt template version · sampling params · verifier result + version judge model + score · dedup cluster · decontam check version · licence / terms of the generator
Generate → verify → judge → dedup → decontaminate → select. Every sample carries a provenance record, so a bad generator version or prompt template can be traced and removed.

Component walkthrough

Failure modes and trade-offs

Interview angle

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.

Go deeper

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

Events (Kafka) impressions, clicks, item create/update LLM feature job new/changed items → batch LLM (tags, scores, embedding), ver vN Streaming aggregations user counts, recency features (Flink / Spark streaming) Offline store (Iceberg / Hive) feature values with event_time + feature version; full history for point-in-time joins Online store (Redis / Cassandra / DynamoDB) latest value per entity; item features also pushed to in-memory cache in rankers Training pipeline labels (from logs) ⋈ features AS OF event_time → train → eval → registry Ranking service same feature definitions; logs features served at request time (log-and-wait)
One feature definition, two stores. The offline store keeps history for point-in-time-correct training joins. The online store keeps the latest values. Logging the features actually served gives the cleanest training data.

Key concepts

Common mistake

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.

Go deeper

Part 3 · Cross-cutting concerns

Reproducibility

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

LeverWhere it appliesCaveat
Spot / preemptibleBatch inference, evals, data processing, small experimentsNeeds small idempotent work units. Rarely used for large synchronous training
Raise utilizationShared clusters (Case 3): scavenger tier, idle reaping, backfillPreemption costs must be tracked
Right-size hardwareSmall models and encoders on cheaper GPUs; decode on memory-bandwidth-rich partsBenchmark tokens per dollar on your workload
Raise MFUTraining: parallelism tuning, FP8, kernel fusion (A3, A7)Engineering time vs GPU savings
Raise goodputLarge runs: faster restart, async checkpoints (Case 2)Most valuable at the largest scale
Cache and reusePrefix caching in serving, stage caching in data pipelines, eval result cachingCache keys must include every version that matters
Reduce workSmaller proxy models for data ablations, tiered evals, early stopping of bad runsProxy 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:

The training-to-serving handoff

training checkpoint (sharded, optimizer state, trainer layout) │ 1. consolidate + strip optimizer state → weights only (≈ 2 B/param in bf16) │ 2. convert to serving format (safetensors, HF layout or engine-specific), attach tokenizer + chat template │ 3. optional: quantize (FP8 / INT4 / NVFP4) with calibration data → re-run evals on the quantized artifact │ 4. eval gates: capability suite, safety evals, regression vs current prod, latency/throughput benchmark │ 5. register: model version → artifacts + eval report + lineage (run, data snapshots) + approvals │ 6. deploy: shadow → canary 1% → 5% → 25% → 100%, gated on SLOs + online quality signals ▼ production alias moved; previous version kept warm for rollback
Go deeper

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:

  1. Exact dedup: group by URL and by content hash.
  2. 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.
  3. 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.

  1. 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.
  2. 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.
  3. Confirm: cross-check the suspects against DCGM, XID errors and NIC counters.
  4. 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:

  1. Label quality: inter-annotator agreement over time, gold-item accuracy, signs of LLM-assisted annotation, and drift in the guidelines.
  2. 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.
  3. Coverage: new prompts may be piling into categories that are already saturated.
  4. 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).