When do you add replicas instead of raising the tensor-parallel degree?
answer
- one factorization, three factors
- smallest group that meets the constraint
- duplicate weights are stolen cache
- only one of them moves per-token latency
- failure domain size is part of the choice
basics
~20 sDefault to replicas whenever the model fits on one GPU: independent copies add throughput linearly with no cross-GPU collectives and no shared failure domain. Shard only when the model does not fit, when per-token latency must drop, or when one replica needs a bigger KV-cache pool.
solid answer
~50 sWith a fixed GPU budget you are choosing between many small serving units and few large ones. **Replicas** win on throughput per dollar: each copy runs alone, there is no all-reduce in the decode loop, scaling is close to linear, and one crashed GPU takes down 1/N of capacity instead of everything. **Shards** (a higher tensor-parallel degree) win in three cases. First, necessity — the weights plus a usable KV cache simply do not fit on one card. Second, latency — decode is memory-bandwidth-bound, so splitting weight reads across ranks lowers time-per-output-token in a way replicas never can. Third, cache capacity: N replicas store N copies of the weights, while a TP=N group stores one copy split N ways, leaving far more memory for KV cache and so supporting longer contexts and larger batches per replica. My rule: take the smallest tensor-parallel degree that meets the memory and latency SLO, and spend every remaining GPU on replicas — with the whole group kept inside one NVLink domain.
go deeper
Understand the two ways to use extra GPUs: run more copies of the model, or split one model across them. Know that copies are the simpler default when the model fits.
Explain the tradeoff mechanically — collective overhead and sub-linear scaling on the shard side, duplicated weights and unchanged per-token latency on the replica side.
Pick a concrete layout for a real fleet and justify it against the binding SLO, including KV-cache headroom per replica, node placement, and what a single GPU failure costs in each option.
Own it as capacity strategy: define the factorization policy for the fleet, whether to split traffic into latency and throughput pools, and how the choice shows up in cost per million tokens and in the blast radius of a hardware failure.
## The decision, framed properly Given G GPUs, a serving layout is a factorization: TP degree x PP degree x replica count = G. The interview question is which factor to grow first, and the defensible answer is *the smallest shard group that satisfies your hard constraints, then replicas for everything else*. ## The case for replicas Replicas are embarrassingly parallel. Each one loads the whole model and serves requests alone, so there is no collective on the decode path, no synchronization stall, and scaling is close to linear in aggregate throughput. Operationally they are also far kinder: a replica is an independent failure domain, so a bad GPU, a stuck NCCL collective or an OOM kills 1/N of capacity rather than the whole serving unit. They roll one at a time during a deploy, they can be scheduled onto whatever nodes have room, and they can even be heterogeneous. For any model that fits comfortably on a single card — most 7B–13B checkpoints, and larger ones once quantized — replicas are the correct default and it is not close. ## The three reasons to shard anyway **1. It does not fit.** The binding constraint is not just weights: you need weights plus a *useful* KV cache plus activation workspace. A checkpoint that technically fits on one GPU but leaves room for two concurrent 8k sequences is not really served. This is the most common and least interesting reason to shard. **2. Latency.** Token-by-token decode reads the weights once per step, so it is bound by memory bandwidth. Tensor parallelism divides those reads across ranks and multiplies aggregate bandwidth, which shortens the step. Replicas cannot do this — a request lives on one GPU either way. If your SLO is written in time-per-output-token and you are already at the bandwidth wall, a higher TP degree is the lever, and it is the reason a latency-tier deployment may run TP=4 on a model that fits on one card. **3. Cache capacity per replica.** This one is regularly missed. Running two replicas on two GPUs stores the weights twice; running TP=2 stores them once, split. The memory you did not spend on a duplicate weight copy becomes KV cache, so a sharded replica supports longer contexts and deeper batches. For a 32k-context workload that difference can be the whole ballgame. ## What sharding costs you Collectives in the inner loop (see the interconnect discussion), sub-linear scaling, and a single fused failure domain — the ranks start together, stall together and die together, and a hung collective usually looks like a wedged server rather than a clean error. Sharding also constrains scheduling: TP=8 means eight GPUs on one node with a good fabric, which is much harder to place than eight independent pods. And the degree is not freely chosen — it must divide the model's head counts, and for compiled engines it may be fixed at build time. ## Working the numbers out loud Suppose 8 GPUs and a model whose weights occupy about 40% of one card. Layout A: 8 replicas, TP=1 — maximum aggregate throughput, each replica has ~60% of a card for KV cache, and any single GPU failure costs 12.5% of capacity. Layout B: 4 replicas at TP=2 — each replica now stores the weights once across two cards, so roughly 80% of two cards is free for KV cache (much longer contexts, bigger batches), per-token latency drops, but you pay per-layer collectives and lose 25% of capacity per failed GPU. Layout C: 1 replica at TP=8 — lowest latency, largest single cache pool, worst throughput per dollar, and the whole service dies with one card. Choose by which SLO is binding: interactive chat under a strict TPOT target argues for B or C, a bulk-scoring job argues for A. ## A note on the knobs The degrees are set per engine and are not portable: vLLM takes `--tensor-parallel-size` (and `--data-parallel-size` for replicas managed by the same launcher), TGI's `--num-shard` is its tensor-parallel degree, and TensorRT-LLM bakes the layout into the engine at build time so changing it requires a rebuild. The reasoning above is engine-agnostic; only the spellings differ. ## Weak answers to avoid "More GPUs on the model is always faster" ignores collective overhead and throughput per dollar. "Replicas are always cheaper" ignores that a duplicated weight copy is memory stolen from the KV cache, and that replicas cannot move time-per-output-token at all.
- Give a case where sharding is right even though the model fits on a single GPU.Two. A strict time-per-output-token SLO: decode is bandwidth-bound, so tensor parallelism splits the per-token weight reads and shortens the step, which replicas cannot do. And a long-context workload: replicas duplicate the weights on every card, while a tensor-parallel group stores one copy split across them, freeing a large amount of memory for KV cache and so admitting far more concurrent 32k sessions per replica.
- How does the failure domain differ between the two layouts?A replica is independent — losing one GPU removes exactly one replica's worth of capacity, and the rest keep serving. A tensor-parallel group is a single unit: all ranks initialize, synchronize and terminate together, so one failed or hung rank takes the whole replica down, often as a stalled collective rather than a clean crash. That argues for health checks with timeouts and for keeping the group no larger than necessary.
- Does mixing degrees across a fleet ever make sense?Yes — a common pattern is two pools behind a router: a small high-TP pool for latency-sensitive interactive traffic and a larger TP=1 replica pool for bulk or batch work. It gives each traffic class the layout its SLO actually needs instead of forcing one compromise, at the cost of running and monitoring two configurations and needing routing logic that knows which class a request belongs to.
saying these in an interview costs you the question
- Assuming more GPUs on one model always means more throughput
- Forgetting that replicas duplicate the weights in every GPU's memory
- Claiming replicas can lower per-token decode latency
- Ignoring that a shard group fails as one unit
- Choosing a tensor-parallel degree without checking node and fabric boundaries