Skip to content

Parallelism

Deployment decisions are often framed as configuration choices such as DP2TP2 versus TP4, but those labels are only shorthand for deeper systems trade-offs. Some forms of parallelism primarily solve a throughput problem, some solve a memory-fit problem, and some introduce communication overhead large enough to erase their theoretical benefit. Parallelism becomes much easier to reason about once those underlying problems are separated.

Data Parallelism (DP)

Data Parallelism (DP) is the most conventional and intuitive form of distributed execution in machine learning systems. Its central idea is simple: the same model is instantiated on multiple devices, while the input workload is partitioned across those instances. Each worker executes an identical computation graph on a different subset of data. In this sense, it changes only the mapping of computation to hardware.

When one model replica is already memory-feasible but cannot absorb enough request load, the natural way to increase system capacity is to add more replicas.

Because the parameters are duplicated rather than partitioned, DP does not reduce the memory footprint on any individual device. This property defines its practical scope: DP is effective when the model already fits within the memory budget of a single GPU, but it is insufficient when the model itself exceeds per-device capacity. Its primary value therefore lies in scaling throughput and serving concurrency through replication.

In training, DP is commonly used to increase effective batch size and improve hardware utilization. Each worker processes a shard of the mini-batch, computes local gradients, and then participates in synchronization so that all model copies remain numerically consistent after the optimization step. In inference, the same replication principle still applies, but without gradient exchange. Instead, each worker handles an independent subset of requests or token batches. Under inference workloads, DP should therefore be understood mainly as a mechanism for traffic distribution rather than model partitioning.

In hybrid deployments, DP often coexists with Tensor Parallelism (TP) and Expert Parallelism (EP). Under such configurations, a DP replica should not be interpreted too literally as a single GPU. Rather, one DP unit may itself be a multi-GPU group whose internal execution is organized by TP or EP. From this perspective, DP defines the outer structure of workload replication, while TP and EP determine how computation is carried out within each replica.

The communication behavior associated with DP depends fundamentally on whether the system is performing training or inference.

During distributed training, each DP worker computes gradients from its local mini-batch shard. Since every worker maintains the same parameter set, these gradients must be synchronized after each step to preserve consistency across replicas. The standard collective primitive for this purpose is all-reduce over gradients. In some implementations, the same effect is realized through reduce-scatter followed by all-gather, but the objective remains unchanged: to aggregate gradient contributions from all DP ranks and produce identical parameter updates on every worker.

By contrast, inference does not involve backward propagation or optimizer steps. Each DP rank serves an independent portion of the request stream, so DP itself introduces no mandatory per-layer collective communication. Under inference, it functions primarily as a scheduling and traffic-partitioning mechanism rather than a synchronization-intensive execution strategy.

That said, a DP-based inference system is not entirely communication-free. The runtime still requires system-level coordination for request dispatch, load balancing, admission control, and result collection. In practice, such coordination often relies on communication patterns such as all-gather to exchange metadata, state, or scheduling information across ranks. These operations belong to the serving framework rather than the model's numerical execution.

In hybrid MoE deployments, DP usually contributes limited communication overhead during inference, while the parallel mechanisms inside each DP unit are often more communication-intensive. Tensor Parallelism typically relies on collectives such as all-reduce or all-gather within a TP group, whereas Expert Parallelism introduces token routing across devices that host different experts. Accordingly, communication analysis in a hybrid system should separate the outer DP dimension from the inner model-parallel dimensions: DP governs workload distribution, while the dominant communication cost usually comes from TP or EP.

Tensor Parallelism (TP)

TP partitions the computation of a single layer across multiple devices by sharding the layer's weight tensors, rather than replicating the full layer on every rank. In large transformer models, this is typically applied to the projection matrices in attention and MLP blocks. The core objective is to reduce the per-device memory footprint and enable execution of layers whose parameters or activations would otherwise exceed the capacity of a single accelerator.

TP becomes interesting when one accelerator cannot comfortably host the model or when reducing per-device weight footprint creates more room for activations and KV cache inside a single replica.

Consider a linear transformation

\[ Y = X W, \]

where \(X \in \mathbb{R}^{B \times H}\), \(W \in \mathbb{R}^{H \times O}\), and \(Y \in \mathbb{R}^{B \times O}\). Under TP, \(W\) is divided across a process group of \(p\) ranks. Two common decompositions are column-wise sharding and row-wise sharding.

In column-wise TP, the output dimension \(O\) is partitioned, so each rank stores

\[ W_i \in \mathbb{R}^{H \times O/p} \]

Each rank computes a partial output

\[ Y_i = X W_i \in \mathbb{R}^{B \times O/p} \]

which corresponds to a slice of the final output tensor. This pattern is natural for the first projection in transformer feed-forward blocks, because the partial outputs can often remain distributed until a later synchronization point.

In row-wise TP, the input dimension \(H\) is partitioned, so each rank stores

\[ W_i \in \mathbb{R}^{H/p \times O} \]

The input \(X\) must then be partitioned consistently, and each rank computes a partial contribution to the full output:

\[ Y_i = X_i W_i \in \mathbb{R}^{B \times O} \]

Since these are additive contributions to the same output tensor, the ranks must sum them to recover \(Y\). This makes row-wise TP naturally coupled with a reduction step.

For transformer layers, TP is usually designed so that adjacent projections alternate between communication-light and communication-heavy layouts. A standard example is the two-layer MLP block. The first linear layer expands the hidden state and is often implemented with column-wise TP, while the second projects back to the hidden dimension using row-wise TP. This pairing avoids unnecessary materialization of full intermediate activations on every rank and keeps most of the computation local until the boundary where the distributed partial results must be merged.

The main advantage of TP is that it scales model width without requiring full parameter replication. However, its efficiency depends strongly on the balance between local matrix multiplication and inter-rank communication. As TP degree increases, per-rank GEMM sizes become smaller and communication becomes a larger fraction of the step time. For this reason, TP is usually effective only within a tightly coupled device group, such as GPUs connected by high-bandwidth intra-node interconnects.

The communication pattern in TP is determined by how tensors are sharded and whether the local computation produces disjoint output slices or partial sums. In practice, TP relies primarily on three collective patterns: all-gather, reduce-scatter, and all-reduce. Their role is not incidental; they define the synchronization boundaries of tensor-parallel execution.

In a column-wise partitioned linear layer, each rank produces a distinct shard of the output features. If the next operator can consume the activation in sharded form, no immediate communication is required. This is one of the main reasons column-wise sharding is attractive. However, when a subsequent operation expects the full activation tensor, the ranks must perform an all-gather to assemble

\[ Y = [Y_0, Y_1, \ldots, Y_{p-1}] \]

Thus, all-gather is associated with reconstructing a feature dimension that was split across ranks.

In a row-wise partitioned linear layer, each rank computes only a partial contribution to the same output tensor. The global result is

\[ Y = \sum_{i=0}^{p-1} Y_i. \]

This requires a summation across the TP group. Conceptually, this is an all-reduce if every rank needs the full output. In optimized implementations, the communication may instead be fused with downstream sharding requirements and expressed as a reduce-scatter, which both sums the partial results and redistributes the output into shards for the next stage.

At the systems level, the cost of TP communication is shaped by message size, collective algorithm, process-group topology, and overlap with computation. Since TP collectives are invoked at fine granularity inside individual layers, their latency sensitivity is high. This is why TP is typically constrained to a small number of nearby devices: the communication volume may be moderate, but the synchronization frequency is high. Once the TP group spans slower links, collective latency can dominate the runtime and erase the gains from parameter sharding. A practical interpretation is that TP converts a single large matrix multiplication into several smaller local GEMMs plus structured collectives. Its performance therefore depends on whether the reduction in memory pressure and per-rank compute fits outweighs the repeated cost of all-gather, reduce-scatter, and all-reduce across the tensor-parallel group.

Inference Interpretation

For serving workloads, DP and TP usually play different roles.

  • DP primarily scales throughput by creating more replicas that can serve independent requests.
  • TP primarily improves model fit and per-replica memory budget by sharding parameters, but it introduces intra-layer collectives.
  • Increasing DP duplicates model weights and may fragment prefix-cache locality across replicas.
  • Increasing TP reduces per-device weight memory and can enlarge the effective KV budget of one replica, but too much TP may turn communication into the bottleneck.

This is why inference deployment decisions are often phrased as a trade-off between replica-level concurrency and per-replica memory headroom.

Case Study: GLM-5.2 with TP8 and EP8

The generic description of TP becomes more concrete in a sparse MoE model such as GLM-5.2. The mapping below reflects the current vLLM implementation of GlmMoeDsaForCausalLM, which reuses the DeepSeek-V2/V3-style MLA and sparse-attention execution path. It describes a single replica configured with TP8 and with expert parallelism enabled across the same eight ranks.

All eight TP ranks process the same request and the same token batch. TP therefore differs from request-level parallelism: it partitions selected model tensors and operators inside one forward pass, while every request still occupies the entire TP group.

Component Placement under TP8 Main communication
Query heads and MLA B projections Split across eight ranks Local attention followed by output reduction
Attention output projection Row-sharded TP all-reduce
Dense MLP gate/up projection Column-sharded Usually none at the expansion boundary
Dense MLP down projection Row-sharded TP all-reduce
Token embedding and LM head Vocabulary-sharded Gather/reduction around logits as required
RMSNorm, RoPE, and residual operations Replicated None or fused with neighboring operations
MLA A projections Replicated None
DSA indexer Largely replicated Index-selection-specific synchronization
Routed MoE experts with EP enabled Experts partitioned across EP ranks Token all-to-all and result exchange
MLA and indexer KV state Full token range retained per TP rank Not sequence-sharded by TP

MLA Attention

Let the model have \(H\) query heads and let the TP degree be \(p=8\). Each rank computes

\[ H_{\mathrm{local}} = \frac{H}{8} \]

query heads. In vLLM, q_b_proj and kv_b_proj are column-parallel projections: their output dimensions are partitioned so that each rank produces the Q/K/V features for its local heads. The attention kernel then operates on those local heads. The output projection is row-parallel, and its partial contributions are summed across the TP group with an all-reduce before the hidden state proceeds to the next block.

MLA also contains low-rank A projections before the head-wise B projections. The current implementation represents q_a_proj and kv_a_proj_with_mqa as replicated linear layers. Every rank therefore stores the same A-projection weights and repeats this part of the computation. TP8 reduces the head-wise B projections and the local attention work, but it does not divide every attention operation by eight.

Dense MLP Layers

Dense MLP layers use the usual paired TP layout. The gate and up projections are fused into a column-parallel linear layer, which partitions the intermediate dimension across ranks. The down projection is row-parallel and combines the local partial outputs with a TP reduction. The intermediate activation can remain sharded between the two projections, avoiding an unnecessary all-gather at the expansion boundary.

This layout reduces the per-rank MLP parameter footprint and GEMM size. Its scaling efficiency falls as TP grows because each local GEMM becomes narrower while the down-projection reduction remains on the critical path.

DSA Indexer and KV State

The sparse-attention indexer is an important exception to the head-wise TP pattern. In the current implementation, its query projection is replicated and the fused key/weight projection explicitly disables TP. Each rank therefore performs largely the same indexer projection and top-\(k\) selection work.

The same distinction applies to KV memory. MLA stores a compressed latent vector for every cached token. Local query heads on every TP rank need access to the complete historical token range, so each rank keeps the compressed MLA state for that range. The indexer cache is similarly maintained over the full sequence. TP8 does not turn one 160k-token KV history into eight independent 20k-token shards.

Consequently, TP and context parallelism solve different problems. TP partitions heads and projection matrices, while DCP/CP partitions the sequence or KV dimension. A TP8 deployment can leave more HBM for KV cache by reducing the per-rank weight footprint, but the eight GPUs do not pool their KV capacity into an eight-times-longer context merely because TP is enabled.

MoE Layers with Expert Parallelism

When expert parallelism is enabled, routed experts follow a different placement rule from attention and dense MLP layers. The router is replicated so that ranks make consistent routing decisions, while the physical routed experts are distributed across the EP group. Tokens are sent to the ranks that own their selected experts, processed by local expert GEMMs, and returned to their original execution ranks.

For TP8 + EP8, the resulting execution structure is therefore:

  • attention projections and heads use TP8;
  • dense MLP layers use TP8;
  • routed experts are partitioned with EP8;
  • the router, MLA A projections, and most indexer work are replicated;
  • shared experts normally retain a tensor-parallel path unless a backend fuses them into the MoE implementation.

This hybrid structure introduces two communication patterns inside the same transformer layer. Attention and dense projections invoke TP collectives, whereas routed MoE computation invokes token all-to-all exchange. Increasing TP or EP changes only the portions assigned to that parallel dimension; replicated operators and communication boundaries remain and prevent single-request latency from scaling linearly with the number of GPUs.

Serving Implications

The placement above leads to several operational consequences:

  • One request consumes all eight GPUs; TP8 does not provide eight independent request slots.
  • Attention-head and large projection work becomes smaller per rank, but replicated MLA and indexer work limits compute scaling.
  • TP collectives are paid at fine granularity during every decode step, so low-concurrency decode often scales less efficiently than prefill.
  • EP reduces routed-expert parameter replication, but token routing adds all-to-all traffic and may suffer from expert-load imbalance.
  • MLA KV capacity remains a per-rank constraint. Higher active-context capacity requires more free HBM per rank, a lower-precision KV representation, sequence/context parallelism, or an offloading design that supports active KV state.

AWQ or other weight quantization changes the storage and kernel used for the sharded matrices, but it does not fundamentally change this parallel placement.

8-GPU Inference Example

The clearest way to compare DP and TP is to hold the hardware fixed and change only how the eight GPUs are grouped into replicas. Consider one 8-GPU server with fast intra-node interconnect. The common options are:

  • DP8TP1: 8 independent single-GPU replicas
  • DP4TP2: 4 replicas, each spanning 2 GPUs
  • DP2TP4: 2 replicas, each spanning 4 GPUs
  • DP1TP8: 1 replica spanning all 8 GPUs

This comparison is only meaningful when all of these layouts are memory-feasible for the target model and context window. If the model or KV state does not fit under TP1, then high-DP layouts are not real options.

The important distinction is that DP spends GPUs on more independent replicas, while TP spends GPUs on making one replica larger. As a result, the best choice depends strongly on whether the workload is constrained by single-request latency, replica memory headroom, or multi-request concurrency.

One subtle point matters a great deal in long-context serving: TP does not only change the compute pattern of one replica. By reducing per-GPU weight footprint inside that replica, it can also leave more memory available for KV cache. That larger KV budget may allow one replica to admit more long requests at the same time, avoid preemption, and form a larger stable batch. In other words, TP can sometimes improve throughput not because one decode step becomes intrinsically faster, but because the replica can sustain enough active work to keep the GPUs busy.

For this reason, inference tuning on 8 GPUs is usually balancing three separate questions at once:

  • How many independent replicas should the node expose?
  • How much KV headroom does each replica need?
  • Is the current workload limited by communication, by memory pressure, or by insufficient active work?

Prefill-heavy workload: 40k + 1k

This workload is dominated by ingesting a very long prompt. Prefill is expensive because attention and projection work scale over a large input sequence before decode even becomes the main issue.

  • Under low concurrency, TP often becomes more attractive because multiple GPUs can cooperate on the same long prefill. DP8TP1 gives many replicas, but if only one or two requests are active, most replicas sit idle while one GPU bears the full prompt cost.
  • Under high concurrency, DP becomes more competitive again because many long prefills can proceed in parallel on different replicas. If each request already fits comfortably, DP8TP1 or DP4TP2 often achieves higher aggregate throughput than very large TP groups.
  • The main reason to keep some TP in this regime is usually not decode speed but memory headroom: very long prompts consume KV capacity quickly, and moderate TP can create a larger feasible per-request context budget.
  • This memory effect can change concurrency as well. With ultra-long prompts, TP1 may force each replica to admit only a very small number of active requests in order to avoid preemption. A moderate TP degree can enlarge the per-replica KV budget enough that the replica can keep a larger prefill batch resident and reach higher real throughput.

In practice, prefill-heavy serving often favors a middle ground such as DP4TP2 when one-GPU replicas are too tight on memory but TP4 or TP8 would sacrifice too much replica count.

Decode-heavy workload: 1k + 10k

This workload is dominated by long autoregressive generation after a relatively short prompt. The system spends most of its time stepping through decode, where each step is small and repeated many times.

  • Under high concurrency, DP is usually the strongest choice if the model fits. DP8TP1 maximizes the number of independent decoders and avoids the per-layer TP collectives that would be paid on every decode step.
  • Under low concurrency, very large TP groups rarely help much unless they are needed for memory. The decode step for one request is narrow, so spreading it across many GPUs can increase synchronization cost faster than it reduces local compute time.
  • If output lengths are large enough to make KV pressure the dominant bottleneck, moderate TP can still win indirectly by enlarging the feasible KV budget of one replica and reducing preemption risk. That is a memory effect, not a pure compute-speed effect.
  • This distinction matters in practice. A decode-heavy service may look compute-light at the per-step level, but if each replica is constantly KV-constrained, it may never accumulate enough active sequences to saturate the device. In that case, adding some TP can increase stable decode concurrency even though it also adds communication.

This is why decode-heavy online serving often prefers as much DP as memory allows. Once the model already fits and KV is stable, extra TP tends to trade throughput for communication.

Balanced workload: 10k + 10k

When prefill and decode are both substantial, neither pure-DP nor pure-TP logic is sufficient. The deployment must absorb long contexts while also surviving long output tails.

  • Under small or moderate concurrency, hybrid layouts such as DP4TP2 are often a robust compromise. They preserve several independent replicas while giving each replica more memory headroom than TP1.
  • Under large concurrency, the decision depends on what fails first. If DP8TP1 starts to preempt or must severely limit admitted requests because KV is too tight, then adding TP can improve real throughput even though it increases communication. If DP8TP1 remains stable, it often keeps the highest aggregate throughput.
  • DP1TP8 is usually justified only when one replica truly needs the entire node, either because the model is too large or because the target context budget is extreme. Otherwise it gives up too much outer concurrency.

This is often the regime where the distinction between logical concurrency and effective concurrency matters most. DP8TP1 may expose more replicas on paper, but if each replica must operate with a very small admitted batch to stay within its KV budget, the system may still fail to keep the node compute-bound. A smaller number of larger TP replicas can sometimes sustain more useful work in flight.

Trade-off Summary

For inference on one 8-GPU node, the qualitative pattern is usually:

  • More DP helps when the workload has enough independent requests to keep many replicas busy.
  • More TP helps when a single replica is too memory-constrained or when one request is so heavy that multiple GPUs should cooperate on it.
  • More TP can also help when ultra-long requests make each TP1 replica too KV-constrained to admit enough active work. In that case, TP is buying concurrency through memory headroom.
  • Prefill-heavy, low-concurrency traffic pushes the system toward some TP.
  • Decode-heavy, high-concurrency traffic pushes the system toward more DP.
  • Balanced long-context traffic often lands in the middle, where DP4TP2 or DP2TP4 is easier to stabilize than either extreme.

The most important operational question is not "is DP better than TP" in the abstract. It is "what is the first wall this workload hits on this hardware: request-level latency, per-replica KV capacity, or cluster-level concurrency?" DP and TP are two different ways of moving that wall.