Model Training Infrastructure and Distributed Training Questions
Scaling model training across hardware and time. Covers GPU/accelerator considerations, data and model parallelism, distributed and large-scale training, experiment tracking and training infrastructure, and the training-versus-inference compute tradeoff. Focuses on the systems and resource decisions that make large-model training feasible.
Explain mixed-precision training: define fp32, fp16 and bfloat16, describe benefits for throughput and memory, and list common numerical pitfalls and mitigations (loss scaling, master fp32 weights, selective casts). Also describe which hardware features (e.g., Tensor Cores) and frameworks make mixed precision safe and efficient.
Sample Answer
Direct answer
Mixed-precision training runs the forward and backward passes in a lower-precision format (fp16 or bfloat16) for speed and memory savings, while keeping a master copy of the weights in fp32 for the optimizer update, since accumulating many small fp16 gradient updates directly into fp16 weights loses precision that matters over the course of training.
Structured elaboration
- fp32: the standard single-precision format (8 exponent bits, 23 mantissa bits), wide dynamic range and good precision but twice the memory and compute cost of fp16/bf16.
- fp16: half precision (5 exponent bits, 10 mantissa bits); much narrower dynamic range than fp32, which is why small gradient values can underflow to zero and large activation/gradient values can overflow to infinity during training without countermeasures.
- bfloat16: also half precision but keeps fp32's 8 exponent bits (same dynamic range as fp32) while sacrificing more mantissa precision (7 bits); this trades some precision for avoiding fp16's overflow/underflow problems, which is why bf16 has become the more common default on hardware that supports it (avoids needing loss scaling in most cases).
- Throughput and memory benefits: half-precision tensors are half the memory footprint and modern GPU tensor cores execute half-precision matrix multiplies at roughly 2x (or more) the throughput of fp32, so both compute and memory bandwidth benefit.
- Common numerical pitfalls and mitigations: fp16's narrow dynamic range means small gradients can underflow to exactly zero, silently stalling learning for those parameters; loss scaling (multiplying the loss by a scale factor before backward, then dividing gradients by that same factor before the optimizer step) shifts small gradient values into fp16's representable range, preventing underflow. Dynamic loss scaling adjusts the scale factor automatically, increasing it when no overflow is observed for a while and decreasing it sharply if an overflow (inf/nan gradient) occurs.
Worked example
A gradient value of 1×10−8 underflows to zero in fp16 (whose smallest positive normal value is around 6×10−5). Scaling the loss by, say, 1024 before backward makes that same underlying gradient appear as 1.024×10−5 after the scale is applied through the chain rule, still within fp16's representable range, and after backward the optimizer divides the accumulated gradient by 1024 again before applying the update, recovering the correct (unscaled) magnitude for the actual parameter update.
Selective casts and frameworks
Selective casting is the mechanism that makes 'mixed' precision actually mixed rather than all-or-nothing: rather than manually annotating every operation, frameworks provide an automatic-cast context (PyTorch's torch.autocast/torch.cuda.amp.autocast, or NVIDIA's older Apex amp) that runs numerically robust, compute-heavy ops (matrix multiplies, convolutions) in fp16/bf16 for the throughput win, while automatically keeping numerically sensitive ops (softmax, loss computation, some reductions and normalization statistics) in fp32, since those operations are more prone to precision-related instability at low bit-width. This selective-cast policy is what lets a model get fp16/bf16's throughput on the operations that benefit most without manually rewriting the model in mixed types. In terms of frameworks and hardware: PyTorch's torch.cuda.amp/torch.autocast and TensorFlow's tf.keras.mixed_precision API are the standard software layers, and both rely on NVIDIA Tensor Cores (available from the Volta generation onward, with native bf16 support from Ampere onward) to execute the low-precision matrix multiplies at higher throughput than fp32 CUDA cores, which is the hardware feature that actually makes the speedup real rather than just theoretical.
Trade-offs & pitfalls
bf16 avoids the underflow/overflow problem that motivates loss scaling in the first place, at the cost of coarser precision per value (fewer mantissa bits), which can matter for numerically sensitive operations (some normalization statistics, certain loss computations) that are often deliberately kept in fp32 even in an otherwise-bf16 training run.
Explain synchronous versus asynchronous stochastic gradient descent in a distributed data-parallel setup. Discuss convergence guarantees, staleness, and scenarios where asynchronous updates are attractive despite potential instability.
Sample Answer
Direct answer
Synchronous SGD has every worker compute a gradient against the same, current parameter values and waits for all workers before applying a single combined update, giving convergence behavior equivalent to (or very close to) single-machine SGD at a larger effective batch size; asynchronous SGD lets each worker push its gradient and pull fresh parameters independently, without waiting for others, trading some workers computing gradients against slightly outdated ("stale") parameters for higher hardware utilization.
Structured elaboration
- Synchronous: every worker's gradient this step is computed against identical parameter values (the state after the previous step's update); once all gradients arrive, they're averaged and applied as one update, after which every worker again has identical, up-to-date parameters. Convergence guarantees closely mirror standard SGD's, since the process is mathematically equivalent to computing a gradient over a larger effective batch (the concatenation of every worker's mini-batch).
- Asynchronous: a worker pulls current parameters, computes a gradient, and pushes it back independently of other workers' progress; by the time its push arrives, the server's parameters may have already been updated by other workers' pushes in the meantime, meaning the pushed gradient was computed against parameters that are now "stale" (out of date) relative to the current server state.
- Staleness and its effect: the degree of staleness (how many other updates happened between a worker's pull and its push) tends to grow with more workers and with heterogeneous worker speeds (a slow worker's gradient becomes more stale the longer it takes to compute); staleness biases the effective update direction, since it's technically a gradient of an earlier point on the loss surface being applied to a later point, which can slow or, in extreme cases, destabilize convergence if unbounded.
- When each is chosen: synchronous is the default for most modern large-scale training (predictable convergence behavior, well-supported by AllReduce-based collectives) provided stragglers are managed; asynchronous is chosen specifically when worker heterogeneity or unreliability is severe enough that waiting for the slowest worker every step would be prohibitively wasteful, accepting some convergence-quality cost in exchange for higher aggregate hardware utilization.
Worked example
With 8 workers, one of which is consistently 3x slower than the others (a straggler): synchronous training's every-step wall-clock time is bounded by that slowest worker, wasting the other 7 workers' idle time waiting each step; asynchronous training lets the 7 faster workers keep contributing updates continuously without waiting, at the cost of the slow worker's occasional contributions being noticeably stale (computed against parameters several updates out of date) by the time they arrive.
Trade-offs & pitfalls
Bounded-staleness schemes (allowing async updates but capping how stale any single contribution is allowed to be before it's rejected or down-weighted) are a common middle ground, retaining most of asynchronous training's utilization benefit while limiting the worst-case convergence-bias risk that fully unbounded asynchrony carries.
Compare AllReduce-based collective communication and a parameter server architecture for distributed gradient aggregation. Explain their basic mechanics, typical frameworks that implement them, and one scenario where one approach outperforms the other.
Sample Answer
Direct answer
AllReduce-based training has every worker communicate directly with every other worker (typically via a ring or tree topology) to jointly compute a fully-reduced result, with no single coordinating node. A parameter server (PS) architecture instead has dedicated server nodes that receive gradients from workers, apply the update, and push back the new parameters, so workers never talk to each other directly.
Structured elaboration
- Communication pattern: AllReduce is peer-to-peer and symmetric: every worker's outbound and inbound traffic is roughly equal. PS is a star/hub pattern: server nodes see traffic proportional to the number of workers, so they can become a bandwidth bottleneck as worker count grows unless sharded across many servers.
- Framework fit: AllReduce is the default in most modern deep learning frameworks (PyTorch DDP, Horovod using NCCL) because it needs no extra server infrastructure and scales bandwidth per-worker rather than concentrating it. PS architectures (originally popularized for very large sparse models) are still common when parameters are enormous and sparsely accessed, since each worker only needs to pull the parameter shards relevant to its mini-batch. Classic named implementations include the original Parameter Server design (Li et al., 2014, used inside systems like Google's DistBelief-era infrastructure) and, in more recent open-source form, TensorFlow's
tf.distribute.experimental.ParameterServerStrategyor a hand-rolled PS built ontorch.distributed.rpc; these give PS a concrete framework grounding the same way PyTorch DDP/Horovod do for AllReduce. - Fault tolerance: PS can tolerate a slow or failed worker more gracefully in asynchronous mode (workers push/pull independently), while synchronous AllReduce requires every worker to reach the collective call, so one straggler blocks everyone.
- Staleness: AllReduce is naturally synchronous (fully consistent gradients each step); PS supports asynchronous updates, trading gradient staleness for higher worker utilization.
Worked example
A 500M-parameter dense model trained on 32 GPUs: AllReduce moves each worker's full gradient (~2GB in fp32) through the ring at ~2x that size in total link traffic per worker, independent of the number of parameter servers you would otherwise need to provision. A PS setup for the same model would need server-side bandwidth sufficient for 32 workers pushing and pulling ~2GB each per step; without sharding the parameters across multiple servers, a single PS node's NIC becomes the bottleneck well before the GPUs do.
Trade-offs & pitfalls
For dense models where every worker touches every parameter, AllReduce is almost always the better default today. PS earns its complexity back mainly for extremely large, sparsely-accessed parameter sets like recommendation-system embedding tables, where each worker only needs a small slice of the total parameters per step.
You're planning on-prem GPU procurement for large-scale training. Draft a plan covering GPU model selection, rack and rack-power sizing (kW per rack), cooling requirements, UPS, networking (RDMA, 100GbE), procurement timelines, vendor support, capacity planning for target utilization, and how these physical choices affect training throughput, single-job latency, and total cost of ownership.
Sample Answer
Direct answer
Planning on-prem GPU procurement for large-scale training requires a plan spanning GPU model selection (matched to the actual workload's compute/memory needs), physical infrastructure sizing (rack space and power, since dense GPU nodes draw far more power per rack than typical enterprise servers), and cooling requirements sufficient for that power density, since underestimating any one of these three typically becomes the binding constraint regardless of how well the other two are planned.
Structured elaboration
- GPU model selection: match the specific GPU generation/model to the actual workload requirements (memory capacity per GPU for the target model sizes, precision support needed like native bf16/fp8, and interconnect capability like NVLink generation), rather than defaulting to "the newest/most powerful available," since procurement lead time and cost scale with how cutting-edge the chosen hardware is, and a slightly older generation may be substantially cheaper and faster to procure while still meeting the actual workload's needs.
- Rack and rack-power sizing: dense GPU servers (e.g. 8-GPU nodes with high-end data-center GPUs) can draw substantially more power per rack unit than typical enterprise compute hardware, often requiring specific high-density rack power provisioning (many kW per rack) well beyond what a typical data center's existing rack power allocation provides; this needs to be sized explicitly against the specific GPU model and node density chosen, not assumed to fit within existing, non-GPU-optimized rack power budgets.
- Cooling requirements: power drawn becomes heat that must be removed; at the power densities modern GPU racks require, some data centers need liquid cooling or other high-density cooling solutions beyond standard air cooling, which is both a cost and a facility-capability question that needs to be resolved (does the target facility support this, or does a facility upgrade/different facility choice become necessary) before hardware even arrives.
- Networking infrastructure: beyond the GPUs themselves, the inter-node networking fabric (InfiniBand or high-speed Ethernet, matched to the target training scale's communication needs) needs its own procurement and installation planning, often on a similarly long lead time to the GPUs themselves, and needs to be sized for the target cluster's actual communication pattern (all-reduce-heavy synchronous training needs particularly high bandwidth and low latency between nodes). Concretely this typically means RDMA-capable fabric (InfiniBand, or RoCE over 100GbE or faster Ethernet) rather than standard TCP/IP networking, since RDMA's ability to transfer data between GPUs across nodes without routing through the host CPU is what makes multi-node all-reduce fast enough not to become the dominant bottleneck at scale.
- UPS (uninterruptible power supply): sized to bridge the gap until backup generators engage, or to allow a controlled checkpoint-and-shutdown of in-flight training jobs, since an ungraceful power loss mid-training risks losing the last unsaved checkpoint's worth of progress and can corrupt a checkpoint file that was mid-write at the moment power dropped; UPS capacity needs to be sized against the same power-density numbers used for rack and cooling planning, not treated as a separate, smaller-scale concern.
- Vendor support: negotiate support/SLA contracts with a bounded hardware-replacement turnaround time, since in a large synchronous training job a single failed GPU or failed NIC can stall the entire job until it's replaced or worked around, and on-site spare-parts agreements for the most failure-prone components meaningfully reduce that downtime risk compared to a standard mail-in RMA process.
- Capacity planning for target utilization: provision against a realistic target utilization (commonly in the 70-85% range) rather than assuming the cluster runs at 100% of nameplate capacity continuously, since job-scheduling gaps, planned maintenance windows, and failed-node downtime all reduce achievable utilization below the theoretical maximum; sizing the cluster as if it will hit 100% utilization leads to a plan that can't actually deliver the training throughput it was sized for.
- Lead time and phased procurement: GPU hardware (especially cutting-edge models) can have long procurement lead times (many months, sometimes longer during periods of high industry demand), which needs to be planned well ahead of the target training start date, potentially with a phased procurement approach (securing an initial tranche of capacity while the remainder is still in the supply pipeline) rather than assuming all hardware arrives simultaneously on a fixed date.
Worked example
Planning a 64-GPU on-prem cluster (8 nodes of 8 GPUs each): GPU model selection settles on a specific generation balancing memory-per-GPU against cost and availability lead time; rack power sizing calculates each 8-GPU node's power draw (summing GPU TDP, host CPU/memory/storage draw, and typical power-supply overhead) and provisions rack power circuits with headroom above that calculated draw; cooling capacity is validated against the facility's existing cooling infrastructure, with a liquid-cooling retrofit budgeted if the calculated heat output exceeds what existing air cooling can handle; and procurement is initiated with enough lead time (informed by current vendor-quoted lead times, not optimistic assumptions) before the planned training start date, with network fabric procurement running on a parallel timeline.
Trade-offs & pitfalls
The most common on-prem GPU procurement planning mistake is treating GPU hardware acquisition as the only real constraint and underestimating power/cooling infrastructure needs, which can become the actual bottleneck (a facility physically unable to support the power density of the newly-arrived GPU hardware) well after the (harder-to-quickly-fix) hardware procurement itself is complete; power and cooling capacity should be validated and, if necessary, upgraded on a timeline that keeps pace with hardware procurement, not treated as an afterthought to be solved once hardware arrives. Each of these choices also has a direct, traceable effect on throughput, latency, and total cost of ownership: insufficient rack power or cooling headroom forces GPUs to thermal-throttle, directly reducing achieved training throughput below the hardware's rated capability regardless of how well everything else was planned; insufficient network bandwidth caps achievable scaling efficiency as GPU count grows (per the earlier scaling-efficiency discussion), which shows up as worse single-job latency at scale; and TCO must account for ongoing power, cooling, and support-contract operating costs, not just GPU sticker price, since a cheaper GPU with worse power efficiency or a weaker support contract can end up costing more over its operating lifetime than a pricier, better-supported alternative.
Design (hard): Create a cost- and carbon-aware ML training scheduler for a hybrid cluster (on-prem + multiple cloud regions). The scheduler should minimize monetary cost and carbon emissions while meeting job deadlines, respecting data locality, and offering preemption options (spot instances). Describe the input signals, objective function, constraints, and how to estimate carbon intensity per region.
Sample Answer
Direct answer
A cost- and carbon-aware scheduler for a hybrid on-prem and multi-cloud training cluster needs two coupled objective functions (minimize dollar cost, minimize carbon emissions) resolved through a configurable weighting or hard constraint, real-time or forecasted signals for both electricity carbon-intensity and cloud spot pricing across candidate locations, and a concrete methodology for actually estimating carbon intensity per region rather than treating it as a given input.
Structured elaboration
- How to estimate carbon intensity per region: the practical approach is to consume a third-party grid carbon-intensity feed (e.g. WattTime, Electricity Maps, or a cloud provider's own published regional carbon-intensity API) that reports grams of CO2 per kWh for a given grid region, typically as both a current/real-time figure and a day-ahead forecast derived from the region's generation-mix forecast (how much of the grid's expected supply is renewable, gas, coal, nuclear at each hour). Two distinct metrics matter here and should not be conflated: AVERAGE emissions intensity (the grid's overall generation mix at that hour) versus MARGINAL emissions intensity (the emissions of the specific power plant that would ramp up or down in response to this workload's incremental demand); for a scheduling decision that shifts load in time or place, marginal intensity is the more correct signal (it reflects the actual emissions impact of the decision), while average intensity is more commonly available and used as a practical proxy when marginal data isn't accessible for a given region. For on-prem sites without a third-party feed covering that specific grid, a fallback is estimating intensity from the site's utility-reported annual generation mix (a coarser, non-time-varying estimate) combined with published regional grid-mix data as a sanity check.
- Cost signal sourcing: spot/preemptible pricing varies by region, instance type, and time.
- Joint optimization approach: express the scheduling decision as an optimization problem, e.g. minimize dollar cost subject to a maximum carbon budget, rather than a single blended score.
- Temporal and spatial flexibility as the actual lever: for workloads with flexibility in when or where they run, the scheduler can shift work to lower-carbon time windows or lower-carbon/lower-cost regions.
- On-prem versus cloud trade-off: on-prem capacity has a largely fixed carbon footprint; cloud capacity offers more flexibility to chase better carbon/cost signals elsewhere.
Worked example
A training job with no hard deadline queries a carbon-intensity API for its three candidate cloud regions and finds region A's day-ahead forecast shows a trough of roughly 120 gCO2/kWh overnight (high wind generation) versus a current 420 gCO2/kWh in region B; scheduling the job to start during region A's forecasted trough rather than immediately in region B cuts its emissions by roughly 3.5x for a modest delay, using the day-ahead average-intensity forecast as the practical signal since marginal-intensity data wasn't available for region A's grid operator.
Trade-offs & pitfalls
Blending cost and carbon into a single weighted score without an explicit, documented rationale tends to produce a policy whose trade-off is opaque. Relying on average rather than marginal carbon intensity (often the only practically available signal) is a known approximation that should be documented as such, since it can misestimate the true emissions impact of shifting a specific workload, particularly in grids where the marginal generator differs substantially from the average mix (e.g. a grid with high average renewable share but a gas-plant marginal generator at the margin).
Unlock Full Question Bank
Get access to all 17 Model Training Infrastructure and Distributed Training interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.