Below the API

Data Parallelism

Advanced Intermediate 45 min Difficulty 2/5

Prerequisites 01


1. What is it?#

Run N complete copies of the model on N GPUs (or N groups of GPUs). Each request goes to one replica. There is no communication between replicas.

       ┌─────────┐  ┌─────────┐  ┌─────────┐  ┌─────────┐
       │ GPU 0   │  │ GPU 1   │  │ GPU 2   │  │ GPU 3   │
       │ full    │  │ full    │  │ full    │  │ full    │
       │ model   │  │ model   │  │ model   │  │ model   │
       │ own KV  │  │ own KV  │  │ own KV  │  │ own KV  │
       └────▲────┘  └────▲────┘  └────▲────┘  └────▲────┘
            │            │            │            │
            └────────────┴──── LOAD BALANCER ──────┘

This is the simplest and most efficient form of parallelism, and for throughput scaling it is almost always the right answer.


2. Why it exists#

Because most inference workloads need throughput, not lower per-request latency, and throughput parallelizes perfectly across independent requests.

1 replica:  1,000 tokens/sec
4 replicas: 4,000 tokens/sec        (linear!)

No AllReduce, no synchronization, no pipeline bubbles. The scaling efficiency is ~100%.

Contrast with TP=4, which gives perhaps 3.4x throughput and adds communication on every layer.


3. Simple analogy#

Four checkout lanes versus four cashiers sharing one till.

Four independent lanes: perfect scaling, no coordination, four customers served at once.

Four cashiers cooperating on one transaction: faster for that customer, but they spend time coordinating and can’t serve four customers at once.

Data parallelism is four lanes. Tensor parallelism is four cashiers on one transaction.


4. Tiny example#

# 4 independent replicas, each on one GPU
CUDA_VISIBLE_DEVICES=0 vllm serve model --port 8000 &
CUDA_VISIBLE_DEVICES=1 vllm serve model --port 8001 &
CUDA_VISIBLE_DEVICES=2 vllm serve model --port 8002 &
CUDA_VISIBLE_DEVICES=3 vllm serve model --port 8003 &

# Load balance across them (nginx, envoy, or your own router)

Or, if the engine supports it natively:

vllm serve model --data-parallel-size 4      # newer vLLM versions

That’s the entire implementation. Compare to tensor parallelism, which requires NCCL initialization, weight sharding, collective operations in the forward pass, and careful topology management.


5. Technical explanation#

What each replica has#

Full model weights            (duplicated N times — the cost)
Its own KV cache pool          (independent)
Its own scheduler and queue    (independent)
Its own CUDA graphs            (independent)

The cost is memory duplication. N replicas need N × weight_bytes. For a 70B FP16 model, 4 replicas = 564 GB of weights. If your GPUs are 80 GB, that’s fine (141 GB per replica needs TP=2 anyway). For an 8B model it’s 64 GB across 4 GPUs — trivially affordable.

The unit of replication#

You replicate whatever fits and meets latency:

8B model:    replicate 1 GPU     → DP across single GPUs
70B model:   replicate TP=4      → DP=2 × TP=4 = 8 GPUs
405B model:  replicate TP=8      → DP=N × TP=8

The composition DP × TP is how essentially every production deployment is structured: TP within a node (over NVLink) for capacity and latency, DP across nodes for throughput.

Routing matters more than it does with TP#

With DP, the load balancer’s choices determine everything:

Round robin      → poor (Section VIII.06); request costs vary 1000x
Least-KV-loaded  → good
Prefix-aware     → best for shared-prefix workloads

A critical consequence: each replica has its OWN prefix cache. Round-robin routing across 8 replicas gives a 1/8 prefix hit rate. Prefix-aware routing gives near-1.0. This is often the largest single factor in a DP deployment’s efficiency.

What DP does NOT give you#

✗ Lower per-request latency. Each request runs on one replica at that
  replica's speed.
✗ Ability to run a model that doesn't fit. Each replica needs the whole model.
✗ A larger effective KV pool for a single request.

6. Under the hood — the memory economics#

Configuration            GPUs   Weight memory   KV memory     Max concurrency
1 × (8B, 1 GPU)          1      16 GB           60 GB         111 @ 4k
4 × (8B, 1 GPU) DP=4     4      64 GB           240 GB        444 @ 4k
1 × (8B, TP=4)           4      16 GB (4 GB/GPU) 300 GB       555 @ 4k

Interesting: TP=4 gives MORE total KV capacity than DP=4, because the weights aren’t duplicated. So if KV capacity is your constraint rather than throughput, TP has an advantage.

But: TP=4 has one shared queue and one scheduler, so a single long request affects everyone; with DP=4, three replicas are unaffected. Isolation is a real benefit of DP.


7. Performance#

                         Throughput   Per-request latency   Efficiency
1 GPU                    1.0x         baseline              100%
DP=4                     4.0x         same                  ~100%
TP=4 (NVLink)            3.4x         4x better             85%
DP=2 × TP=2              3.6x         2x better             90%

DP=2 × TP=2 is often the sweet spot when you need both some latency improvement and good throughput scaling: you get most of TP’s latency benefit at a smaller TP degree (where scaling efficiency is highest) plus DP’s near-perfect throughput scaling.


8. Production implications#

  • Scale throughput with DP. Always. It’s free efficiency.
  • Use the smallest TP degree that meets your fit and latency requirements, then DP on top.
  • Invest in the router. With DP, routing quality determines prefix hit rate and load balance, and both matter a lot (Section VIII.06).
  • DP gives fault isolation. One replica crashing affects 1/N of traffic. A TP group crashing affects all of it.
  • DP gives deployment flexibility. You can roll replicas one at a time; a TP group must be updated atomically.
  • Watch for the memory duplication cost with large models — it’s what pushes you toward TP.

9. Common mistakes#

Using TP when DP would do. The single most common distributed-inference error.

Round-robin routing across DP replicas with prefix caching enabled. You built a cache and then guaranteed misses.

Assuming DP scaling is automatic. It’s only linear if your router balances well.

Forgetting that each replica has its own KV pool. A request can only use one replica’s memory, so DP doesn’t help a single request needing 100 GB of KV.

Not exploiting the isolation. DP lets you canary one replica, drain one replica, and isolate noisy tenants. Use it.


10. Hands-on exercise#

A. Measure DP scaling. Run 1, 2, and 4 independent replicas of a small model behind a load balancer. Measure aggregate throughput. Is it linear? If not, what limited it (the LB, the client, the network)?

B. DP vs TP, head to head. With 4 GPUs, compare DP=4 and TP=4 on: aggregate throughput, p50 latency, p99 latency, and total KV capacity. Build the comparison table.

C. The routing effect. With 4 DP replicas and prefix caching enabled, measure the aggregate prefix hit rate with round-robin and with session-affinity routing. Quantify the throughput difference.

D. Composition. With 8 GPUs, compare DP=8×TP=1, DP=4×TP=2, DP=2×TP=4, DP=1×TP=8 for a model that fits on one GPU. Plot throughput and p95 latency for each. Where’s the sweet spot for your SLO?

E. Fault isolation. Kill one replica in a DP=4 setup and measure the impact. Then kill one GPU in a TP=4 setup. Compare the blast radius.


11. Interview questions#

  1. What is data parallelism and what is its communication cost?
  2. Why is DP preferred over TP for throughput scaling?
  3. What does DP not give you?
  4. Why does TP=4 give more total KV capacity than DP=4?
  5. Why does routing quality matter more with DP than with TP?
  6. What are the fault-isolation differences between DP and TP?
  7. With 8 GPUs and a model that fits on one, how would you configure it and why?

12. Further reading#

  • [ESTABLISHED] Pope et al., “Efficiently Scaling Transformer Inference” (2022)
  • [REFERENCE] vLLM data-parallel serving documentation
  • Next: 03 — Tensor parallelism

↑↓ navigate ↵ open