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.
Compare GPUs and TPUs for model training. Explain workload characteristics where GPUs are a better fit and where TPUs offer advantages. Discuss programming model differences (TensorFlow vs XLA), precision support (FP16, BF16), and practical considerations when selecting instance types for cloud training.
Sample Answer
Direct answer
GPUs are the default choice for most training workloads today because of their broad software ecosystem and flexibility across model architectures, while TPUs offer an edge specifically for very large, highly-regular workloads (especially transformers at scale) where their systolic-array matrix-multiply design and purpose-built high-bandwidth interconnect (ICI) pay off most.
Structured elaboration
- Workload characteristics favoring GPUs: irregular or research-stage architectures (custom ops, dynamic control flow, frequent architecture changes), smaller-scale training, and any workload needing the widest possible library/framework support (most cutting-edge research code targets GPU first).
- Workload characteristics favoring TPUs: very large, well-established architectures (transformers) trained at massive scale, where TPU pods' purpose-built interconnect and systolic-array matrix units are specifically tuned for the dense matrix multiplies dominating transformer compute, often at a favorable cost-per-FLOP for sustained large training runs.
- Programming model differences: GPU programming (CUDA, and the PyTorch/TensorFlow ecosystems built on it) is eager-execution-friendly and broadly flexible; TPUs traditionally require (or strongly favor) XLA-compiled, graph-based execution, which handles static, regular computation graphs excellently but is less forgiving of dynamic shapes or data-dependent control flow, and has historically meant a narrower (though growing) software ecosystem outside Google's own frameworks (JAX, TensorFlow).
- Availability and lock-in: GPUs are available across essentially every cloud provider and on-prem; TPUs are Google Cloud-specific, which is a real practical constraint for teams not already committed to that ecosystem.
Worked example
A research team iterating rapidly on novel architectures with custom CUDA kernels and dynamic control flow would find GPUs' flexibility essential; the same team, once they've settled on a stable, well-established transformer architecture and are scaling a single well-defined training run to thousands of accelerators, might specifically evaluate TPU pods for that production-scale run given TPUs' interconnect and systolic-array advantages at that specific workload shape, while continuing GPU-based experimentation for anything still in flux.
Precision support and instance-type selection
- Precision support: both platforms support bf16, but their history differs. GPUs (from the Volta generation onward for fp16, Ampere onward for native bf16 Tensor Core support) support fp16 well but fp16's narrow dynamic range typically requires loss scaling; TPUs were designed around bf16 from early generations specifically to avoid the loss-scaling machinery fp16 needs, so bf16-native training is the more idiomatic default on TPU, while GPU training commonly still supports both fp16 (with loss scaling) and bf16 depending on the GPU generation and framework defaults.
- Practical considerations for instance-type selection: cloud GPU instances (e.g. A100/H100-based VM families) are available in flexible single-GPU-to-multi-GPU configurations across essentially every major cloud provider, so a team can start small and scale incrementally; TPU pods are typically procured in fixed pod-slice sizes (e.g. v4/v5 slices of a given chip count) on Google Cloud specifically, which means committing to a minimum scale upfront and to a single cloud provider, a real practical constraint when a team wants provider flexibility or needs to start at very small scale before committing to a large training run.
Trade-offs & pitfalls
The GPU-versus-TPU choice is as much about ecosystem and organizational fit (existing tooling, cloud provider commitments, team CUDA expertise) as it is about a workload's raw performance characteristics; a team defaulting to GPUs purely out of ecosystem familiarity, even for a workload that would benefit from TPU's advantages, is a common and often reasonable trade-off, not a mistake, given the switching cost.
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.
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).
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.
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.
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.