PidokuInfra

Why Distribute (and When Not To)

Advanced Intermediate 1h Difficulty 3/5 Topic 01 of 12

Prerequisites V.06, I.08


1. What is it?#

Using more than one GPU for a single model. There are exactly three reasons to do it, and one very common bad reason.

GOOD REASON 1 — CAPACITY: the model doesn't fit on one GPU.
GOOD REASON 2 — LATENCY:  one GPU's bandwidth gives unacceptable ITL.
GOOD REASON 3 — KV SPACE: you need more memory for concurrent sequences.

BAD REASON   — "more GPUs must be faster." Often false. See file 12.

2. Why the distinction matters#

Because the three good reasons lead to different parallelism strategies, and choosing the wrong one wastes money.

REASON                    STRATEGY               WHY
Model doesn't fit         TP, or PP              splits the weights
Latency too high          TP                     more aggregate bandwidth per token
Need more KV capacity     TP (splits KV too),    
                          or more replicas       
Need more throughput      DATA PARALLELISM       replicate; no communication at all
                          (more replicas)

If you only need throughput, do not use tensor parallelism. Replicating the model across GPUs gives linear throughput scaling with zero communication overhead. TP gives sublinear scaling and communication cost. People reach for TP by reflex; usually they wanted replicas.


3. Simple analogy#

Moving a piano.

If the piano fits in your van, one van is optimal. Two vans doesn’t make the delivery faster; it makes it more expensive.

If the piano doesn’t fit, you need two vans and a plan for coordinating them — and now you have coordination overhead you didn’t have before.

If you have ten pianos to deliver, ten vans each carrying one piano is perfect parallelism: no coordination at all. That’s data parallelism, and it’s the best kind when it applies.


4. The decision procedure#

STEP 1: Does the model fit on one GPU, with room for KV cache?
        weights + 2 GB overhead + (KV for your target concurrency) ≤ GPU memory

        YES → do you meet your latency SLO on one GPU?
              (decode floor = weight_bytes / bandwidth)
              YES → USE ONE GPU. Scale with REPLICAS for throughput. Done.
              NO  → go to step 2

        NO  → go to step 2

STEP 2: How many GPUs are needed just to FIT?
        n_fit = ceil(weights / (gpu_memory - overhead - min_kv))

STEP 3: How many GPUs are needed for LATENCY?
        n_latency = ceil(weight_bytes / (bandwidth × ITL_target))

STEP 4: n = max(n_fit, n_latency), rounded up to a power of 2
        (most TP implementations require the number of KV heads to be divisible)

STEP 5: Check the interconnect (file 08). If these n GPUs are not
        NVLink-connected, reconsider — TP over PCIe is often worse than
        a smaller TP degree plus more replicas.

STEP 6: Scale THROUGHPUT with replicas of this n-GPU unit.

5. Worked examples#

Example A — Llama-3-8B, chat, ITL target 40 ms#

Step 1: weights FP16 = 16.1 GB. On an 80 GB H100: fits with 60 GB for KV. ✓
        decode floor = 16.1/3350 = 4.8 ms ≪ 40 ms ✓
        → ONE GPU. Scale with replicas.

Using TP=2 here would: halve the ITL to 2.4 ms (unnecessary),
                       halve the KV per GPU (no capacity gain — same total),
                       add an AllReduce per layer (~15% overhead),
                       and cost 2x per replica.
→ Strictly worse. Use one GPU.

Example B — Llama-3-70B, ITL target 30 ms#

Step 2: weights FP16 = 141 GB. On 80 GB GPUs:
        n_fit = ceil(141 / (80 - 4 - 10)) = ceil(141/66) = 3 → round to 4

Step 3: n_latency = ceil(141e9 / (3.35e12 × 0.030)) = ceil(141/100.5) = 2

Step 4: n = max(4, 2) = 4 → TP=4

Step 5: Check topology. If GPUs 0-3 are NVLink-connected: good.
        If GPUs 0-3 span two PCIe islands: bad. Use TP=4 within one island,
        or reconsider.

Result: TP=4. ITL ≈ 141/4/3350 = 10.5 ms + comms ≈ 14 ms ✓ (target 30)
        Scale throughput with replicas of the 4-GPU unit.

Example C — Llama-3-405B#

weights FP16 = 810 GB.
n_fit = ceil(810 / 66) = 13 → but TP degrees are constrained by head counts.
405B has 128 attention heads, 8 KV heads.
TP must divide n_kv_heads (8) or use KV replication.
→ TP=8 (one node) gives 101 GB/GPU. Doesn't fit in 80 GB.
→ TP=8 with FP8 quantization: 50.6 GB/GPU. Fits. ✓
→ Or TP=8 × PP=2 across two nodes at FP16.

Practical answer: quantize to FP8 and use TP=8 on one node.
Quantization avoided a multi-node deployment — a large operational simplification.

Example C’s lesson generalizes: quantization is often the cheapest alternative to adding parallelism.


6. The four forms of parallelism#

DATA PARALLELISM (DP)
  Replicate the whole model N times. Requests go to one replica.
  Communication: NONE.
  Scales: throughput linearly.
  Doesn't help: latency, or fitting a model that doesn't fit.

TENSOR PARALLELISM (TP)
  Split each weight matrix across GPUs. Every GPU works on every token.
  Communication: AllReduce per transformer block (2 per layer).
  Scales: latency (more aggregate bandwidth), capacity.
  Costs: communication proportional to hidden size and batch.

PIPELINE PARALLELISM (PP)
  Split by layers. GPU 0 has layers 0-19, GPU 1 has 20-39, etc.
  Communication: activations between stages (small).
  Scales: capacity.
  Costs: pipeline bubbles; doesn't reduce per-token latency.

EXPERT PARALLELISM (EP)
  For MoE: different experts on different GPUs.
  Communication: All-to-All (tokens to their experts and back).
  Scales: capacity for MoE models.
  Costs: All-to-All is expensive and load-imbalanced.

Plus sequence/context parallelism (file 06) for very long context.

These compose: a large deployment might be TP=8 within a node × PP=2 across nodes × DP=6 replicas = 96 GPUs.


7. Performance implications#

Strategy       Latency        Throughput      Efficiency     Comms
1 GPU          baseline       baseline        100%           none
DP × N         same           N×              ~100%          none
TP=N (NVLink)  ~N× better     ~0.85N×         85%            AllReduce/layer
TP=N (PCIe)    ~2× better     ~0.4N×          40%            AllReduce/layer
PP=N           same           ~0.7N×          70%            activations

Read the TP rows carefully. With NVLink, TP=4 gives roughly 3.4x the throughput of one GPU and 4x lower latency — a good trade. Over PCIe, TP=4 gives about 1.6x the throughput of one GPU using 4 GPUs, which is a terrible trade. The interconnect determines whether TP is a good idea.


8. Production implications#

  • Default to the smallest parallelism degree that fits and meets latency. Then scale with replicas.
  • Check nvidia-smi topo -m before choosing a TP degree. Keep TP groups within an NVLink domain.
  • Quantization is often cheaper than parallelism. FP8 halves the model; it may turn a multi-node deployment into a single-node one.
  • TP degree must divide the number of KV heads (or the implementation must replicate KV heads, which wastes memory). A model with 8 KV heads supports TP up to 8 cleanly.
  • Multi-node is a step change in operational complexity (file 10). Avoid it if quantization or a smaller model would do.

9. Common mistakes#

Using TP when you wanted DP. The most common error. If you need throughput, add replicas.

Choosing TP=8 because the node has 8 GPUs. Choose from the requirements, then check topology.

Ignoring the interconnect. TP over PCIe is often worse than not doing TP.

TP degree not dividing the KV head count. Either fails or silently replicates KV.

Going multi-node before trying quantization.

Assuming linear scaling. TP scaling is 70-90% per doubling at best.


10. Hands-on exercise#

A. Run the decision procedure for three models you might deploy, on hardware you have access to. Show all six steps. Compare your answer to what the default configuration would be.

B. Measure TP scaling. For a model that fits on one GPU, measure throughput and ITL at TP=1, 2, 4. Plot both. What’s the scaling efficiency at each degree? Compare to the table in section 7.

C. DP vs TP. With 4 GPUs, compare (i) TP=4 on one model instance, and (ii) 4 independent replicas. Measure aggregate throughput and per-request latency for each, under a realistic load. Which is better for throughput? For latency?

D. The topology check. Run nvidia-smi topo -m on your hardware. Which GPU groups are NVLink-connected? Design your TP grouping accordingly. If you have PCIe-only GPUs, measure TP=2 across them and compare to the NVLink case if available.

E. Quantization vs parallelism. For a model that needs TP=4 at FP16, check whether FP8 lets it run at TP=2. Measure both configurations’ throughput and cost.


11. Interview questions#

  1. What are the three legitimate reasons to use multiple GPUs for one model?
  2. Why is data parallelism preferred when you only need throughput?
  3. Walk me through choosing a TP degree for a 70B model with a 30 ms ITL target.
  4. Why must the TP degree relate to the number of KV heads?
  5. When is quantization a better answer than parallelism?
  6. Why does the interconnect determine whether TP is worthwhile?
  7. What’s the scaling efficiency of TP, and why isn’t it linear?

12. Further reading#

  • [ESTABLISHED] Pope et al., “Efficiently Scaling Transformer Inference” (2022) — the definitive partitioning analysis
  • [ESTABLISHED] Shoeybi et al., “Megatron-LM” (2019)
  • [REFERENCE] vLLM distributed serving documentation
  • Next: 02 — Data parallelism

↑↓ navigate↵ openesc close