An ML systems engineer makes models train faster, serve cheaper and fail less, working between the researchers who choose the architecture and the product engineers who call the model over an API. The questions are concrete: why is this run using 40 percent of the GPU, why did the loss go to NaN, how many users will this replica hold, and why did p99 latency triple on a long document.

This roadmap is the companion to the AI engineer roadmap, which deliberately stops at the API boundary, and a track of the 2026 developer roadmap. Here we go below it, from the hardware up to fleet operations. Each stage has a skill to build, a check that proves you have it, and reading on this site.

Advertisement

The stage map

Every performance question in this field reduces to three resources: arithmetic throughput, memory capacity and bandwidth, and the links between devices. Learn to ask which one is the bottleneck before touching any knob.

Where an ML systems engineer works: one request, one training stepData loaderCPU workers, pinned memoryTraining stepforward, backward, optimizerCollectivesall-reduce / all-gather (NCCL)Checkpoint + registryweights, config, eval reportRouter / gatewayauth, quotas, SLO classSchedulercontinuous batching, prefill vs decodeKV cachepaged blocks, prefix reusedeployGPU: SMs + tensor cores | HBM | NVLink / networkevery box above is ultimately a question about compute, memory or bandwidthTraining is throughput-bound and synchronous; serving is latency-bound and bursty. Both are memory problems first.
The two halves of the work share one substrate. Training moves gradients between GPUs; serving moves requests through a scheduler and a KV cache. Both are bounded by compute, memory and interconnect.
StageSkillYou have it when
1. GPU fundamentalsExecution model, memory hierarchy, rooflineYou can say whether a kernel is compute- or bandwidth-bound and why
2. The training stepForward, backward, optimizer, mixed precisionYou can estimate memory for a model before launching it
3. Scaling outData, tensor, pipeline parallelism; sharding; collectivesYou can explain where every all-reduce and all-gather happens
4. Inference internalsPrefill vs decode, KV cache, quantizationYou can size the KV cache for a workload
5. ServingBatching, scheduling, routing, SLOsYou can load-test a replica and read the latency curve
6. OperationsProfiling, failures, cost, fleet schedulingYou can debug a slow or dead job from its telemetry

Stage 1: how a GPU actually runs your code

A GPU is many streaming multiprocessors (SMs), each running thousands of threads grouped into warps that execute the same instruction together. Tensor cores inside each SM perform small low-precision matrix multiplies far faster than the ordinary arithmetic units. Data lives in a hierarchy: registers and shared memory on the SM are tiny and fast, the L2 cache is shared, and HBM holds the model but is, relative to the arithmetic, slow.

The single most useful mental model is the roofline. Every kernel has an arithmetic intensity: floating-point operations per byte moved from HBM. Large matrix multiplications have high intensity and saturate the tensor cores. Elementwise operations, layer norms and, crucially, decode-time attention have low intensity and are limited by memory bandwidth. This is why kernel fusion and FlashAttention matter: they avoid writing intermediate results to HBM and reading them back.

Precision is part of the hardware story. BF16 is the default training format on current data-centre GPUs. Hopper introduced FP8 matrix math through NVIDIA's Transformer Engine, and Blackwell adds 4-bit floating point (NVFP4) and MXFP8 with fine-grained block scaling. Lower precision doubles or quadruples throughput per byte only if the scaling keeps values in range; the formats are a numerical-methods problem as much as a hardware one.

Check: write a naive CUDA or Triton matrix multiply, profile it, then compare it with the library version and explain the gap in terms of memory traffic.

Read next on this site: GPU architecture overview, the SIMT execution model, the GPU memory hierarchy, tensor cores, HBM, occupancy, kernel fusion, FlashAttention on the GPU, FP8.

Advertisement

Stage 2: the anatomy of a training step

One step runs the forward pass and keeps activations, runs the backward pass to produce gradients, and applies the optimizer. Memory is the first constraint you will hit, and you should be able to estimate it before launching anything. With mixed-precision Adam, each parameter costs about 16 bytes of model and optimizer state before a single activation is stored, so a full fine-tune of a 7-billion-parameter model needs over 100 GiB for states alone and does not fit on one 80 GB GPU without sharding or offload.

# Back-of-envelope GPU memory. Numbers you should be able to do in your head.
GiB = 1024 ** 3

def training_bytes(params, bytes_per_param=16):
    # Mixed-precision Adam (ZeRO paper): bf16 weights (2) + bf16 grads (2)
    # + fp32 master weights, momentum, variance (4 + 4 + 4) = 16 bytes/param.
    # Activations come on top and depend on batch, sequence and checkpointing.
    return params * bytes_per_param

def kv_cache_bytes(layers, kv_heads, head_dim, tokens, dtype_bytes=2):
    # K and V for every layer, every KV head, every token held in the cache.
    return 2 * layers * kv_heads * head_dim * dtype_bytes * tokens

print(f"7B full fine-tune states: {training_bytes(7e9) / GiB:.0f} GiB")
per_tok = kv_cache_bytes(layers=80, kv_heads=8, head_dim=128, tokens=1)
print(f"KV per token (70B-class, GQA): {per_tok / 1024:.0f} KiB")
print(f"32 seqs x 4096 tokens: {kv_cache_bytes(80, 8, 128, 32 * 4096) / GiB:.0f} GiB")

Activations scale with batch size, sequence length and depth. Activation checkpointing recomputes them in the backward pass instead of storing them. Gradient accumulation runs several micro-batches per optimizer step, giving a large effective batch on limited memory.

Numerical stability is the second constraint. Watch the loss and the gradient norm together; a gradient-norm spike usually precedes a loss spike. BF16 has the range of FP32 and rarely needs loss scaling, while FP16 does. A NaN is a bug to root-cause, not a thing to skip: log the step, the batch indices and the checkpoint so the failure can be reproduced.

Read next on this site: training step anatomy, mixed precision, optimizer offload to CPU, full fine-tuning, LoRA and QLoRA.

Stage 3: scaling out across GPUs

Data parallelism gives every GPU a full copy of the model and a different slice of the batch, then averages gradients with an all-reduce. It is the first tool because it is simple and the all-reduce can overlap with the backward pass. When the model or its optimizer state no longer fits, shard: ZeRO and fully sharded data parallelism split optimizer state, gradients and finally parameters across ranks, gathering each layer's weights just before use.

Tensor parallelism splits individual matrix multiplies across GPUs and needs a collective inside every layer, so it belongs on the fast NVLink domain within a node. Pipeline parallelism splits layers into stages and passes activations between them; its cost is the bubble while the pipeline fills and drains. Mixture-of-experts models add expert parallelism and an all-to-all exchange of tokens. Large runs combine these; put the chattiest collective on the fastest link.

The worked loop below is data parallel with gradient accumulation. Note no_sync, which skips the all-reduce on the micro-batches that do not end an optimizer step; forgetting it multiplies network traffic by the accumulation factor.

# torchrun --nproc_per_node=8 train.py
import contextlib, os, torch, torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP

dist.init_process_group("nccl")
rank = int(os.environ["LOCAL_RANK"]); torch.cuda.set_device(rank)
model = DDP(build_model().cuda(), device_ids=[rank])
opt = torch.optim.AdamW(model.parameters(), lr=3e-4, weight_decay=0.1)
ACCUM = 4                                   # micro-batches per optimizer step

opt_step = 0
for step, (x, y) in enumerate(loader):      # DistributedSampler inside loader
    x, y = x.cuda(non_blocking=True), y.cuda(non_blocking=True)
    sync = (step + 1) % ACCUM == 0
    ctx = model.no_sync() if not sync else contextlib.nullcontext()
    with ctx:                                # all-reduce fires only when sync
        with torch.autocast(device_type="cuda", dtype=torch.bfloat16):
            loss = model(x, y) / ACCUM      # scale so accumulated grad is a mean
        loss.backward()
    if sync:
        gnorm = torch.nn.utils.clip_grad_norm_(model.parameters(), 1.0)
        if not torch.isfinite(gnorm):       # check BEFORE the update lands
            raise RuntimeError(f"non-finite grad norm at step {opt_step}")
        opt.step(); opt.zero_grad(set_to_none=True); opt_step += 1
        if rank == 0 and opt_step % 1000 == 0:
            save_checkpoint(model, opt, opt_step)   # weights + optimizer + data position

Check: run the same job on one and on eight GPUs, compute scaling efficiency, and use a profiler timeline to show whether communication overlaps with compute.

Read next on this site: NCCL collectives, overlapping communication with compute, ZeRO sharding, pipeline parallelism, expert parallelism, NVLink, InfiniBand.

Stage 4: inference is a different workload

Generation has two phases. Prefill processes the whole prompt in parallel and is compute-bound, much like training's forward pass. Decode produces one token at a time per sequence, reading every weight and the whole KV cache for each token, and is bandwidth-bound. Time to first token is dominated by prefill; time per output token by decode. Most serving optimisations are ways to keep decode busy.

The KV cache stores the keys and values of every previous token so they are not recomputed. It grows linearly with context length and concurrency, and for long contexts it outgrows the weights. For a 70B-class model with grouped-query attention shaped like Llama 3 70B (80 layers, 8 KV heads of dimension 128), each token costs about 320 KiB in BF16, so 32 sequences of 4,096 tokens need about 40 GiB. Grouped-query and multi-head latent attention exist largely to shrink this number.

Weight quantization shrinks what decode must read and KV-cache quantization raises how many tokens fit; both change outputs, so each needs an evaluation. Speculative decoding lets a cheap draft propose tokens that the large model verifies in one pass.

Read next on this site: the KV cache, KV-cache quantization, multi-head latent attention, quantization calibration, speculative decoding, sampling and decoding.

Stage 5: serving systems and SLOs

A serving engine is a scheduler wrapped around a KV-cache allocator. Continuous batching admits new requests into the running batch at every decode step instead of waiting for a whole batch to finish. PagedAttention, introduced by vLLM, stores the cache in fixed-size blocks mapped through a block table, so memory is allocated on demand and fragmentation is small. vLLM's automatic prefix caching and SGLang's RadixAttention reuse the cached prefix shared by requests with the same system prompt or documents. Chunked prefill splits long prompts so they do not stall everyone else's decode, and prefill-decode disaggregation runs the two phases on separate GPU pools.

The worked example: a team must serve a 70B model to an internal assistant with p95 time to first token under one second. Weights in BF16 are about 130 GiB, so one 80 GB GPU is out and two leave almost no KV headroom. Four GPUs with tensor parallelism inside one NVLink domain give the budget below.

# Capacity sketch for one serving replica: will the KV cache fit the traffic?
hbm_gib         = 4 * 80          # 4 GPUs with tensor parallelism
weights_gib     = 70e9 * 2 / 2**30   # 70B params in bf16 ~ 130 GiB
overhead_gib    = 0.10 * hbm_gib     # activations, CUDA graphs, fragmentation
kv_budget_gib   = hbm_gib - weights_gib - overhead_gib
kv_per_tok_kib  = 320                # from the estimator above
max_tokens      = kv_budget_gib * 2**20 / kv_per_tok_kib
avg_ctx         = 3000               # prompt + generated, measured from logs
print(f"KV budget {kv_budget_gib:.0f} GiB -> ~{max_tokens/avg_ctx:.0f} concurrent sequences")

That gives about 158 GiB of cache and roughly 170 concurrent sequences at the measured average context, an upper bound before scheduler headroom. The team load-tests with a realistic prompt-length distribution, plots latency against offered load, and sets the admission limit below the knee of the curve. If average context triples, capacity drops by two thirds; the options are FP8 weights, KV-cache quantization, more replicas or a per-tier context limit, each with a measured trade-off.

Read next on this site: continuous batching, paged attention, prefix caching at scale, chunked prefill, prefill-decode disaggregation, vLLM in production, LLM serving architecture, SLO-aware scheduling.

Stage 6: operating GPU fleets

At fleet scale, hardware failure is routine. A long training job must checkpoint weights, optimizer state and data-loader position, and must be able to resume on a different set of nodes. Measure goodput, the share of paid GPU time that produced progress kept in a checkpoint, not just utilisation. Slow nodes are worse than dead ones because a synchronous job runs at the speed of its slowest rank; per-rank step timing finds them.

Serving needs canary releases compared on quality as well as latency, a registry that ties weights to their evaluation report, and autoscaling on queue depth and KV utilisation rather than CPU. Tokens per GPU-hour at a given latency target is the cost number that decides hardware and quantization choices.

Read next on this site: GPU profiling, nvidia-smi, Slurm scheduling, MIG partitioning, LLM canary releases, model registries, GPU cost optimisation.

Failure modes and trade-offs

  • Optimising the wrong resource. Tuning kernels while the data loader starves the GPU. Profile first; look for gaps on the GPU timeline.
  • Benchmarks that do not match traffic. A fixed 512-token prompt hides the long-context tail that exhausts the KV cache in production.
  • Silent quality loss. Quantization or a new kernel that passes a speed test and degrades answers. Every performance change gets an evaluation run.
  • Checkpoints that cannot resume. Weights saved without optimizer state or data position restart training from a different point in effect.

The standing trade-offs: throughput against latency, memory against recomputation, precision against quality, and specialised speed against operational simplicity. Record each choice and its measurement in the design doc.

What to do next

  • Profile one matrix multiply and one elementwise kernel and place both on a roofline.
  • Run the memory estimator for a model you use, then compare it with measured peak memory.
  • Train a small model with DDP on two GPUs, with gradient accumulation and no_sync, and measure scaling efficiency.
  • Serve a model with vLLM or SGLang, load-test with a realistic prompt-length distribution, and find the latency knee.
  • Size the KV cache for your own traffic and decide on an admission limit.
  • Quantize the served model to FP8 or INT8 weights and gate the change on an evaluation set.
  • Write a resume-from-checkpoint test and run it by killing a job mid-epoch.
Key takeaway: ML systems engineering is resource accounting under uncertainty. Learn how a GPU spends compute, memory and bandwidth, then how a training step uses them, how parallelism places collectives on links, and how inference splits into compute-bound prefill and bandwidth-bound decode around a KV cache. Serving is a scheduler around that cache, judged by SLOs under real traffic. Estimate before you launch, profile before you tune, and gate every speed-up on an evaluation.