Multi-GPU fine-tuning goes wrong for a reason that has nothing to do with your code. Two cards can solve two different problems, and those problems want opposite things. One problem is that the model will not fit on one card. The other is that it fits, but the run takes too long. Aim the wrong fix at your problem and you lose a week.
There is a third fact sitting under both, and it is not about the cards at all. The moment one job spans two cards, those cards have to send data to each other. That data travels down a real physical path at a real speed. That speed usually decides whether the whole idea was worth it.
By the end of this part you will be able to turn a bandwidth number into milliseconds of waiting, say what DDP, ZeRO and FSDP each do to memory and to traffic, and choose between them from one question about your own machine. Nothing about distributed systems is assumed. The previous part covered preference tuning. This one is about getting any run onto more than one card.
- Part 1. LLM Fine-Tuning Explained: What Actually Changes Inside the Model
- Part 2. Training Memory: The Four Tenants and the 16 Bytes Per Parameter
- Part 3. Activation Memory: Why the Forward Pass Costs More Than the Weights
- Part 4. Number Formats for Training: FP32, FP16, BF16, TF32, FP8 and NF4
- Part 5. Gradients and Optimizers: From SGD to Adam, and Why m and v Cost 8 Bytes
- Part 6. AdamW Explained, Line by Line
- Part 7. Masking in LLM Training: Loss, Causal and Padding
- Part 8. Choosing a Base Model and Building a Fine-Tuning Bench
- Part 9. Supervised Fine-Tuning End to End: The First Real Run
- Part 10. Data Is the Actual Job: Formats, Quality, Splits and Synthetic Data
- Part 11. LoRA Explained: Freeze the Model, Learn a Low-Rank Patch
- Part 12. QLoRA Explained: A 4-Bit Base, NF4, and What It Really Costs
- Part 13. Preference Tuning: RLHF, DPO, and the Verifiable-Reward Branch
- Part 14. Multi-GPU Fine-Tuning: DDP, ZeRO, FSDP and the Wire Between the Cards (you are here)
- Part 15. Evaluating a Fine-Tune: The Ladder, the Delta and Six Silent Failures
Tell apart the two problems a second card can solve
Start with what a second GPU physically is.
A GPU is a card that plugs into a slot on the motherboard. It brings its own pool of fast memory, called VRAM, and its own arithmetic units. Put two cards in one machine and you have two pools and two sets of units. They are separate. Nothing moves from one pool into the other unless some code copies it there, one byte at a time.
That single fact splits this whole subject in two.
Problem A, it will not fit. A full fine-tune costs 16 bytes for every parameter in the model. Qwen2.5-7B has 7.62 billion parameters, so 7.62 billion times 16 comes to about 122 GB. That figure is static state only, meaning weights, gradients and optimizer state. Activations are extra and are counted separately. A 32 GB card is nowhere near.
Problem B, it is too slow. The job does fit, but one epoch, meaning one full pass over your training data, takes nine hours. What you want here is throughput: more training examples finished per second.
Two problems, two families of answer, and the words for them are worth getting straight before anything else.
- Replicate fixes B. Every card gets a complete copy of the model. Each card works on a different slice of the data. Nothing gets smaller. More work happens at once.
- Shard fixes A. Take one thing, cut it into pieces, and give each card one piece. No card holds the whole thing. Card A keeps the first half of the numbers, card B keeps the second half.
Hold that difference. A copy on each card, against a piece on each card. Two people can each carry a complete phone book, or they can each carry half of one. The pair with half a book each can still find any name between them. Neither of them can do it alone.
Both families pay the same tax. Copies have to be kept identical, and that takes traffic. Pieces have to be fetched when they are needed, and that takes traffic too.
Turn a bandwidth number into milliseconds of waiting
“The cards have to talk to each other” is a phrase that hides the entire problem. Here is what it means physically.
To get a number out of card A and into card B, that number is copied along a physical path between the two slots. There are two kinds of path, and which one you have is decided by the machine, not by the cards.
- NVLink is a dedicated link between the two cards, built for exactly this job. NVIDIA puts it on their datacentre parts.
- PCIe is the general-purpose slot that every card sits in. With no NVLink, bytes going from card A to card B travel up the slot to the CPU end of the board and back down the other slot. It is a detour, and the road is shared with everything else in the machine.
The link between two cards, whichever kind it is, is the interconnect. Its speed is quoted as bandwidth, which is simply how many bytes per second it can carry. GB/s means billions of bytes per second.
Four numbers frame the whole field. An H100 datacentre card gets 900 GB/s of NVLink. The older A100 gets 600. A full PCIe 5.0 slot gets 128. Two cards sharing one board’s lanes get 64.
That last one needs a word of its own. A PCIe slot is a bundle of lanes, and a full slot is 16 of them. Most consumer boards do not have enough lanes to give two cards 16 each, so adding the second card drops both to 8. Half the lanes, half the bandwidth, and you did nothing wrong.
Treat all four as headline figures. They count traffic in both directions at once, and a real transfer only ever gets some fraction of them. What matters here is not the exact value but the gap between the top of that chart and the bottom.
Now make it concrete, because a number in GB/s means nothing until it becomes a number in seconds. One division does it.
seconds = bytes to move / bytes per second
Take a real pile of bytes. A 7B model in bf16, the compact 2-byte number format this series uses throughout, carries one gradient per weight at 2 bytes each. So 7.62 billion times 2 is about 15 GB of gradient. The next section explains why that exact pile has to cross the wire once per training step. For now, just divide.
- On a 900 GB/s link: 15.24 divided by 900 is 0.017 seconds, about 17 milliseconds.
- On a 64 GB/s link: 15.24 divided by 64 is 0.24 seconds, about 240 milliseconds.
Put those next to the actual work. Suppose the arithmetic in one training step takes 300 milliseconds. The exact figure does not matter, the comparison does. On the fast link, 17 milliseconds of waiting is noise. On the slow one, 240 milliseconds nearly doubles the step. Same code, same model, same batch. The wire decided.
What consumer cards have, and what they do not
Recent consumer GPUs have no dedicated card-to-card link at all. The RTX 3090 was the last consumer card to carry NVLink. It was not the only one: Turing’s 2080 and 2080 Ti carried it before that. Everything after the 3090, including the 40-series and the whole 50-series, has none of it. On those machines, every byte between two cards takes the PCIe detour through the CPU.
This does not make multi-GPU training useless. It makes the choice of strategy decisive. You want the strategies that send less, or send less often. You avoid the ones that have to talk inside every layer.
Follow one training step across two cards under DDP
The baseline everybody starts from is distributed data parallel, DDP for short. Both halves of that name earn their place.
- Distributed means the job runs as several programs at once, one per card, kept in step with each other.
- Data parallel means the data gets split between them. The model does not.
Here is one step, in five moves, on two cards with a batch of 16 examples.
- Each card loads a complete copy of the model, plus the gradients and optimizer state that go with it. That is the full 16 bytes per parameter, on each card, twice over.
- The batch is split. Card A takes 8 examples, card B takes the other 8. The slice one card handles in one go is called a micro-batch.
- Each card runs its own forward pass and its own backward pass on its own 8 examples. Each now holds one gradient for every weight in the model. The two sets of numbers disagree, because the two cards saw different examples.
- The cards combine their gradients so that both end up holding the same numbers. This is the all-reduce, and it is the only moment of coordination in the entire step.
- Each card runs its optimizer step on those shared numbers. Identical gradients applied to identical weights give identical new weights. The two copies stay one model.
What an all-reduce actually does
Take a single weight and follow it through move 4.
Card A’s 8 examples produced a gradient of -0.4 for that weight. Card B’s 8 produced +0.2. Neither is wrong. They are answers to different questions.
The all-reduce does two things in order. First it adds them up: -0.4 plus 0.2 is -0.2. Then it hands that total back to both cards, and each divides by the number of cards to get the average, -0.1.
That -0.1 is what one card would have computed from all 16 examples in one go, provided the two halves were the same size. Which is the whole point. Splitting the batch has not changed the maths of training at all.
The name is literal. “Reduce” means many numbers combined into one. “All” means every card gets the answer back, not just one of them. An operation that every card takes part in together is called a collective, and all-reduce is the one that matters most here.
Why must every card get the answer? Because every card applies its own optimizer step to its own copy of the weights. Give two copies different gradients and they drift apart. A few steps later you are training two different models and have no way to say which one is yours. The all-reduce is what keeps them a single model.
Now the size of it. One gradient per weight, 2 bytes each, means the pile crossing the wire is as big as the model. About 3 GB for Qwen2.5-1.5B. About 15 GB for a 7B. Every step, all run long. That is where the 240 milliseconds in the last section came from.
Why it hurts less than it sounds, until it does not
The backward pass does not produce all the gradients at once. It works backwards, so the last layer’s gradients are ready first and the first layer’s are ready last. There is no reason to sit and wait for the whole set. Frameworks group finished gradients into buckets and start sending each bucket while the remaining layers are still computing. Traffic and arithmetic run at the same time.
On a fast link that hides nearly all of the cost. On a slow one the arithmetic finishes first, and the cards sit there waiting for bytes. Overlapping can only hide what the wire is able to carry. This is why the same script is quick on one machine and painful on another.
Which brings up the most common misunderstanding in the field. DDP does nothing for memory. Every card still holds a full copy of everything. A model that needed more than 32 GB on one card still needs more than 32 GB on each of two. What you bought is throughput, and that is all.
Cut the traffic without changing your model
Two levers shrink the bill before you change strategy at all. Both are settings you already know from earlier parts.
Lever one: talk less often
Gradient accumulation means running several micro-batches one after another, adding each one’s gradients into a running total on the card, and only doing the all-reduce and the optimizer step after the last one.
Set gradient_accumulation_steps to 8 and you get one all-reduce per 8 micro-batches instead of 8 of them. The same examples are processed and the same arithmetic is done. An eighth of the traffic. Part 9 introduced accumulation as a way to hold activation memory down, and here it is doing a completely separate second job.
What changes is that the weights now update 8 times less often per epoch. The number to keep an eye on is the effective batch: the per-card micro-batch, times the accumulation steps, times the number of cards. Two cards at micro-batch 4 and accumulation 8 gives an effective batch of 64.
That number is not free to grow. Goyal and colleagues showed that a large batch trains about as well as a small one if you multiply the learning rate by the same factor you multiplied the batch by, and ramp that rate up gradually over the first stretch of training rather than applying it cold. So adding a second card, or raising accumulation, moves a second dial whether you touch it or not.
The rule runs out eventually. You and colleagues found that past a certain batch size plain scaling stops holding, and built an optimizer, LAMB, that sets a separate step size for each layer in order to push batches into the tens of thousands. Fine-tuning rarely goes anywhere near that scale. It is worth knowing the ceiling exists before you assume a bigger effective batch is always safe.
Lever two: send less
Only trainable weights have gradients. Frozen weights, meaning weights the training loop has been told never to change, have none. LoRA freezes the base model and trains a small adapter instead, roughly 1 percent of the parameters. So roughly 1 percent of the gradient bytes need to cross the wire.
For a 7B model, that turns about 15 GB per step into about 150 MB. On the 64 GB/s link, 2.4 milliseconds instead of 240.
| Setup, Qwen2.5-7B on two cards | Bytes crossing the wire per micro-batch | Time at 64 GB/s |
|---|---|---|
| Full fine-tune, DDP, sync every micro-batch | about 15 GB | about 240 ms |
| Full fine-tune, DDP, accumulate 8 micro-batches | about 1.9 GB | about 30 ms |
| LoRA adapter, DDP, sync every micro-batch | about 150 MB | about 2.4 ms |
| LoRA adapter, DDP, accumulate 8 micro-batches | about 19 MB | under 1 ms |
The middle rows are the same 15 GB all-reduce, spread over eight times as much work. LoRA is usually sold as a memory technique. On a bandwidth-starved box it is a traffic technique too, and the two savings stack.
Shard the four tenants, one ZeRO stage at a time
DDP wastes memory in a very visible way. Two cards hold two identical copies of the weights, two identical copies of the gradients, and two identical copies of the optimizer state. Every byte exists twice, and no card is any better off for it.
ZeRO is the scheme that removes that duplication. The name stands for zero redundancy optimizer. Instead of every card holding the full 16 bytes per parameter, each card keeps one slice of the training state and fetches the rest when it needs it. This is sharding, applied to the four tenants Part 2 named: weights, gradients, optimizer state and activations.
It goes in three stages, and it takes the fattest tenant first.
- Stage 1 shards the optimizer state. That is 12 of the 16 bytes, so it is the biggest single win available. Each card keeps optimizer numbers only for its own slice of the weights, and updates only that slice. The updated weights then get shared so every card is current. The volume on the wire works out about the same as DDP’s all-reduce, so stage 1 is nearly free.
- Stage 2 shards the gradients as well. A card that only updates its own slice of the weights only ever needs gradients for that slice. So instead of every card ending up with every gradient, the gradients are added up and handed out in pieces. Again the volume is about the same as before.
- Stage 3 shards the weights themselves. Now no card holds the whole model at any time. Memory drops as far as it can go. This is the stage that costs, because the weights have to be gathered layer by layer during the forward pass and again during the backward pass.
Do the arithmetic on two cards. Start from 2 bytes of weights, 2 of gradients and 12 of optimizer state.
- Stage 1: 2 plus 2 plus 6, so 10 bytes per parameter per card.
- Stage 2: 2 plus 1 plus 6, so 9 bytes per parameter per card.
- Stage 3: 1 plus 1 plus 6, so 8 bytes per parameter per card.
Stage 3 is just 16 divided by the number of cards. That is the useful form to remember, because it scales: on 8 cards it is 2 bytes per parameter, on 16 cards it is 1.
| Stage | What it shards | Bytes per parameter per card, on 2 cards | Qwen2.5-7B static state per card | Traffic against DDP |
|---|---|---|---|---|
| None, plain DDP | nothing, a full copy each | 16 | about 122 GB | the baseline |
| Stage 1 | optimizer state | 10 | about 76 GB | about the same |
| Stage 2 | and gradients | 9 | about 69 GB | about the same |
| Stage 3 | and the weights | 8 | about 61 GB | roughly half as much again |
Every number in that fourth column is static state. Activations are on top of all of them.
Now read the column honestly, because it says something people do not expect. Not one of those figures fits a 32 GB card. Sharding divides, it does not shrink. Two consumer cards cannot full fine-tune a 7B model, no matter which stage you pick.
What would it take? Stage 3 puts 122 GB divided by the card count on each card. Divide by 8 and you get about 15 GB each, which leaves room for activations on a 32 GB card. So a 7B full fine-tune is an eight-card job. On two cards, the answer is LoRA, and on one card it is LoRA as well.
The same arithmetic scales up cleanly. A 70B model at 16 bytes per parameter is 1,120 GB of static state. Put it on 8 datacentre cards of 80 GB each and stage 3 leaves 140 GB per card, which does not fit. Put it on 16 and you get 70 GB per card, which fits an 80 GB card with activations making it tight. That is the arithmetic behind a 70B full fine-tune being a sixteen-card job rather than an eight-card one.
One tenant never appears in any of this. No stage shards activations. Each card computes activations for its own micro-batch only, so they get smaller because each card’s batch is smaller, not because anything was split. Part 3’s levers, checkpointing above all, still apply per card and are still worth pulling.
Two names go with the stages. ZeRO was published together with a library called DeepSpeed, from the same group, whose config files name the stages outright as zero_stage. FSDP is PyTorch’s own implementation of the same idea, and it is the one most fine-tuning stacks reach for now. That is the next section.
Watch one layer’s weights appear and vanish under FSDP
FSDP stands for fully sharded data parallel. Its full sharding mode is stage 3 by another name, and Hugging Face’s own guide states that mapping directly.
Start from what the cards hold when nothing is happening. Card A has half of every layer’s weights. Card B has the other half. Neither card can run the model as it stands. That sounds fatal, and it is not, because a forward pass only needs one layer at a time.
Here is the life of layer 5 on two cards.
- The forward pass reaches layer 5. Card A holds half its weights, card B holds the other half.
- The two cards all-gather layer 5: each sends its slice to the other. Both now hold the layer in full, for a moment.
- Both cards compute layer 5 on their own micro-batch.
- Each card frees the half it does not own. Memory drops straight back down.
- On to layer 6, and the same again.
All-gather is a different collective from all-reduce, and the difference is worth one line each. Gather means everybody ends up holding all the pieces, laid side by side, unchanged. Reduce means everybody ends up holding the pieces added together. Weights want gather, because you need each number exactly as it is. Gradients want reduce, because you want their sum.
The same overlap trick from DDP applies here. While the cards compute layer 5, they start fetching layer 6. That is called prefetching, and it is why a well-configured FSDP run is much faster than the step list above sounds.
The backward pass mirrors it. Each layer is gathered again, its gradients are computed, and then the cards do a reduce-scatter: the gradients are added across cards, and each card is handed only the slice covering the weights it owns. Nobody needs a gradient for a weight they will never update.
So the peak is not the whole model. At any instant a card holds its permanent slice of everything, plus one layer in full. Across 28 layers that is a small fraction of what DDP would have needed.
The price is traffic, and you can count it. A ring all-reduce is really two passes over the data, one to combine and one to hand the result back. FSDP makes three: gather the weights going forward, gather them again going backward, then scatter the gradients. So budget roughly half as much traffic again as DDP, in exchange for the memory.
What a shard unit is
FSDP does not split the model weight by weight. You tell it what a unit is, through a setting called the wrap policy. The standard choice is one transformer block per unit, a transformer block being one full layer of the model, of which Qwen2.5-1.5B has 28. Each block is gathered, computed on and freed as a piece.
| FSDP setting | Equivalent ZeRO stage | What it shards | Reach for it when |
|---|---|---|---|
NO_SHARD |
none, this is plain DDP | nothing, a full copy per card | the model fits and you only want throughput |
SHARD_GRAD_OP |
stage 2 | gradients and optimizer state | you are close to fitting and want lighter traffic |
FULL_SHARD |
stage 3 | weights, gradients and optimizer state | the model will not fit any other way |
The rule that falls out of the table is the same one the stages gave you. Pick the lowest setting that makes the run fit. Do not pay full sharding’s traffic bill if SHARD_GRAD_OP already got you under the card’s limit.
One naming wrinkle to know about
Everything above describes the original FSDP API, the one most training stacks still default to and the one the settings in the table belong to. PyTorch has since rewritten its tutorial around a newer API called FSDP2, built on a function named fully_shard, and now describes the original as deprecated. In FSDP2 a single flag, reshard_after_forward, replaces the settings above: true behaves like FULL_SHARD and false behaves like SHARD_GRAD_OP. Recent Accelerate releases can select it with fsdp_version: 2. Check which one your stack picks before assuming any flag name here still applies.
Set up a real FSDP run in the three places it lives
A working FSDP run is three separate pieces in three different formats. Mixing up which setting lives where is a common first-hour mistake, so each one is labelled below.
Piece one, a shell command. This is typed into a terminal. It is not Python.
accelerate launch --num_processes 2 --multi_gpu train.py
accelerate launchstarts your training script for you instead of running it directly. Accelerate is the Hugging Face library that handles the multi-card plumbing.--num_processes 2starts two copies of the script, one per card. A process is one running program with its own memory, so this really is two programs, not two threads.--multi_gputells Accelerate that those two copies are one job and have to be wired together.train.pyis your ordinary training script, unchanged.
Piece two, a YAML config file. YAML is a plain-text settings format where indentation carries meaning. Normally you generate this by running accelerate config and answering its questions. The FSDP settings are keys nested under fsdp_config, and they are not Python assignments.
compute_environment: LOCAL_MACHINE
distributed_type: FSDP
mixed_precision: bf16
num_processes: 2
fsdp_config:
fsdp_sharding_strategy: FULL_SHARD
fsdp_auto_wrap_policy: TRANSFORMER_BASED_WRAP
fsdp_transformer_layer_cls_to_wrap: Qwen2DecoderLayer
fsdp_backward_prefetch_policy: BACKWARD_PRE
fsdp_state_dict_type: SHARDED_STATE_DICT
distributed_type: FSDPpicks sharded data parallel over plain DDP.mixed_precision: bf16runs the arithmetic in the 2-byte format while a 32-bit copy of each weight is kept for the optimizer.fsdp_sharding_strategy: FULL_SHARDis stage 3. Change it toSHARD_GRAD_OPfor stage 2.fsdp_auto_wrap_policy: TRANSFORMER_BASED_WRAPmakes each transformer block one shard unit.fsdp_transformer_layer_cls_to_wrapnames the class in the model’s code that counts as a block. For Qwen2.5 models it isQwen2DecoderLayer. Get this name wrong and nothing gets wrapped, quietly.fsdp_backward_prefetch_policy: BACKWARD_PREfetches the next layer’s weights before they are needed. This is the overlap described above.fsdp_state_dict_type: SHARDED_STATE_DICTsaves a checkpoint as one file per card, rather than assembling the entire model on one card to write it out.
Piece three, Python. These are the training settings themselves, in the same config object earlier parts used.
from trl import SFTConfig
config = SFTConfig(
output_dir="qwen-fsdp-run",
per_device_train_batch_size=1,
gradient_accumulation_steps=8,
gradient_checkpointing=True,
bf16=True,
optim="paged_adamw_8bit",
max_length=1024,
)
from trl import SFTConfigimports the settings object from TRL, the Hugging Face library that runs supervised fine-tuning.output_diris the folder checkpoints get written to.per_device_train_batch_size=1is per card, not the total. Two cards at 1 means 2 examples in flight.gradient_accumulation_steps=8is lever one. One all-reduce per 8 micro-batches.gradient_checkpointing=Trueis Part 3’s activation lever, and it is still needed here, because no sharding stage touches activations.bf16=Trueasks for the 2-byte number format in the arithmetic.optim="paged_adamw_8bit"is AdamW with its two running averages stored in 1 byte each instead of 4, and able to spill to system memory during a spike. It needs the bitsandbytes package installed.max_length=1024caps how long each training sequence may be, which caps the activation bill.
Pin your library versions before building anything on top of this. Multi-GPU support moves faster than any other area in this series, and flag names in particular have moved between releases.
When the wire is fast on paper and slow in practice
The collectives above are carried out by NCCL, NVIDIA’s library for exactly this. It picks a transport when the job starts, and a handful of environment variables can push it onto a slower path with no error message at all.
NCCL_P2P_DISABLEturns off direct card-to-card transfers, so traffic goes the long way through host memory.NCCL_P2P_LEVELcaps how close two cards must be on the board before they are allowed to talk directly.NCCL_IB_DISABLEturns off InfiniBand, a datacentre network type that a desktop board never had in the first place.
Any of these set wrongly, often inherited from a container image or a copied environment file, costs real speed and says nothing. Set NCCL_DEBUG=INFO before a serious run so the log states which transport was actually chosen. Reports of instability on multi-card consumer machines are common, because that combination is far less tested than the datacentre one. Budget time for tuning, and prefer the most-supported path over hand-rolled code.
Rule tensor and pipeline parallelism in or out
Everything so far splits the batch. Even under full sharding, each card still runs the whole model’s arithmetic on its own examples. Two other approaches split the model itself instead, and both are worth understanding mainly so that you can rule them out on purpose.
Tensor parallelism splits a single matrix multiply. A layer works by multiplying the numbers flowing through it against a grid of weights. Cut that grid down the middle, give card A the left half and card B the right half, and each card produces half the answer. The two halves have to be joined before the next step can run. That means a collective inside every layer, on every pass, in both directions during training. It is the most talkative scheme there is, and it needs a fast link to be worth anything.
Pipeline parallelism splits the layers. Card A holds layers 1 to 14, card B holds 15 to 28. The only thing crossing the wire is the activations handed over at the boundary, which is a tiny amount of traffic. The cost is idle time: while card A works on the first micro-batch, card B has nothing to do, and at the end of the batch it is the other way round. Those gaps are called bubbles. Feeding many small micro-batches through shrinks them and never removes them.
There is a useful asymmetry here that explains something you may have seen. Tensor parallelism is common for serving models on two consumer cards, and it works fine. Two reasons it survives PCIe where training would not. Serving does a single forward pass, with no backward pass, no gradient all-reduce and no optimizer step to synchronise. And serving engines are engineered hard to overlap and minimise what little traffic remains. The part on continuous batching and paged attention covers that engineering. Training the same way would add the same per-layer collectives in the backward direction on every step, which is precisely what PCIe handles worst.
At the largest scale all three axes combine. Tensor parallelism stays inside one node, meaning one machine whose cards share a fast internal link. Pipeline parallelism spans nodes over the slower network between machines. Data parallelism sits on top for throughput. The design keeps the chatty traffic inside the fast domain and sends only cheap boundary activations across the slow links. That is the same principle you apply on a two-card desktop, with one slow wire instead of a hierarchy of fast ones.
Pick a strategy for the hardware you actually have
All of it collapses into one question, asked first: does the job fit on one card?
Work the question honestly. Add the static state, add the activations, then leave 2 to 3 GB of headroom for fragmentation and one unusually long batch.
- It fits. Then the second card is a throughput multiplier. The simplest and often the fastest option is one job per card, which uses no interconnect at all. Use DDP instead when you want one single run to finish sooner.
- It does not fit, and you have not tried LoRA. Try LoRA or QLoRA first, with gradient checkpointing on. For a 7B on a 32 GB card this almost always ends the problem, and it cuts the wire traffic by about ninety-nine percent as a side effect.
- It does not fit, and it genuinely has to be a full fine-tune. Now use FSDP, at the lowest sharding setting that makes it fit, with bf16, gradient checkpointing, and accumulation of at least 8. Expect the gathers to cost real time on PCIe.
- Never tensor-parallel a training run on consumer cards. Keep that one for serving.
The honest headline for two consumer cards: you get roughly twice the throughput and two separate memory pools. You do not get one pool of double the size. Without a dedicated link there is no way to treat two cards as a single big-memory device, so a model too large for one card needs explicit sharding and pays PCIe for it.
For the most common case by far, which is adapter-based fine-tuning of a base model that already fits, the second card is best spent on a parallel experiment or on keeping an inference service alive. That sidesteps the interconnect completely, and it is usually the quickest route to a result.
Whether the run took one card or eight, finishing it is not the same as proving it worked. The final part, on evaluating a fine-tune, is what tells you that.
Key takeaways
- Separate the two problems before anything else. Will not fit needs sharding. Too slow needs replication. Conflating them picks the wrong tool.
- The interconnect is the hidden limit, and it turns into milliseconds with one division: bytes to move divided by bytes per second. A 7B gradient is about 15 GB, which is 17 ms on a 900 GB/s link and 240 ms on a 64 GB/s one.
- Recent consumer cards have no NVLink. The RTX 3090 was the last one to carry it, and the 2080 and 2080 Ti had it before that. Everything since goes over PCIe, through the CPU.
- DDP replicates the model, splits the batch and all-reduces the gradients so every copy applies the same update. It buys throughput and nothing else.
- ZeRO shards one more tenant per stage: optimizer state, then gradients, then weights. Stage 3 costs 16 divided by the card count in bytes per parameter, and sharding divides rather than shrinks, so two 32 GB cards still cannot full fine-tune a 7B.
- FSDP gathers each layer’s weights just in time, computes, frees them and prefetches the next. Budget roughly half as much traffic again as DDP for it.
- Gradient accumulation and LoRA both attack the traffic directly, and they are the two cheapest wins on a slow wire.
You can now
- Say which of the two problems you actually have, and therefore which family of answer to reach for, from “Tell apart the two problems a second card can solve”.
- Convert any interconnect’s bandwidth figure into milliseconds of waiting per step, from “Turn a bandwidth number into milliseconds of waiting”.
- Explain what an all-reduce does to one weight and why every card needs the result, from “Follow one training step across two cards under DDP”.
- Work out per-card memory at any ZeRO stage on any number of cards, and say whether a run will fit, from “Shard the four tenants, one ZeRO stage at a time”.
- Write the three pieces of a working FSDP run in their correct formats, and check that NCCL chose the transport you expected, from “Set up a real FSDP run in the three places it lives”.
- Rule tensor parallelism out for training and in for serving, with the reason, from “Rule tensor and pipeline parallelism in or out”.
Glossary
- 8-bit optimizer
- AdamW with its two running averages stored in 1 byte each instead of 4, using block-wise quantisation. The algorithm and its behaviour are unchanged, and it reclaims most of 8 bytes per parameter. Dettmers et al., 8-bit Optimizers via Block-wise Quantization
- Activation
- Any intermediate value the forward pass produces on its way from input to output. Activations are held in memory because the backward pass needs them to work out the gradients. Activation memory, from birth to death
- Activation memory
- The memory holding the forward pass’s intermediate values until the backward pass consumes them. Its size comes from batch size, sequence length, hidden size and layer count, never from parameter count. Why the forward pass costs more than the weights
- AdamW
- Adam with the weight decay applied straight to the weight instead of folded into the gradient. It is the default optimizer for essentially every LLM fine-tune, and its stored state is 12 of the 16 bytes per parameter. AdamW explained, line by line
- Adapter
- A small set of extra trainable weights added beside a frozen model, so you train the adapter and leave the model alone. A LoRA adapter is tens of megabytes against a multi-gigabyte model copy. LoRA explained
- All-reduce
- The step where every GPU hands over its gradients, they are added up, and the total goes back to all of them. It moves as much data as the model is big, every single step.
- Attention
- The step where each position in the sequence looks at other positions and mixes in whatever it finds useful. It is what lets a model use context instead of reading each token in isolation. Multi-head attention explained
- Backpropagation
- The procedure that computes a gradient for every weight in one sweep backwards through the model, from the loss at the end to the first layer. It works by applying the chain rule one step at a time. The calculus behind backpropagation
- Backward pass
- Running backwards from the loss through the model to produce a gradient for every weight. It is backpropagation in practice, and it needs the activations the forward pass stored. How a neural network learns
- Batch
- A group of examples processed together in one step, so the GPU stays busy and the gradient is averaged over several examples instead of one. Bigger batches give a steadier signal and cost more activation memory. Supervised fine-tuning end to end
- bf16
- A 16-bit number format with 8 exponent bits and 7 mantissa bits, so it reaches as far as fp32 with much coarser steps. It is the training default because it needs no loss scaling. Number formats for training
- Checkpoint
- A saved copy of the model at some point in the run, so you can resume from it, compare it, or ship it. Distinct from gradient checkpointing, which is a memory trick with an unfortunately similar name. Supervised fine-tuning end to end
- Data parallelism
- Putting a full copy of the model on every GPU, splitting the batch between them, and averaging the gradients each step. It buys throughput and does nothing at all for a model that will not fit.
- Effective batch
- The number of examples that actually go into one weight update: the per-device batch times the accumulation steps times the number of GPUs. This is the number that matters for training behaviour, rather than the batch that happens to fit. Batch size, accumulation and the effective batch
- Epoch
- One full pass over the training data. Most instruction fine-tunes need only one to three, and preference tuning usually needs exactly one. Supervised fine-tuning end to end
- Forward pass
- Running data through the model from input to output to get a prediction and a loss. Along the way it produces the activations that the backward pass will need. The complete inference path
- Frozen weights
- Weights marked as not trainable, so they never receive an update. A frozen weight needs no gradient and no optimizer state, which removes 14 of its 16 bytes, though its activations are still stored. LoRA explained
- FSDP
- Fully sharded data parallel, PyTorch’s own implementation of staged sharding. Each card holds a slice of the model at rest and gathers a layer’s full weights just before computing on it, then frees them again. Zhao et al., PyTorch FSDP
- Full fine-tuning
- Updating every weight in the model. It costs 16 bytes per parameter of static state under standard mixed-precision AdamW, produces a complete model copy per task, and forgets the most. Training memory and the 16 bytes per parameter
- Gradient
- One number per weight saying which way to nudge that weight to make the loss smaller, and how steeply the loss responds. Picture the slope under a ball rolling into a valley. How a neural network learns
- Gradient accumulation
- Running several small batches, adding their gradients together, and only updating the weights once at the end. You get the steadier signal of a big batch while holding just one small batch in memory. Batch size, accumulation and the effective batch
- Gradient checkpointing
- Throwing away most stored activations and recomputing them during the backward pass. Peak activation memory drops a long way in exchange for roughly 20 to 30 percent more time. Chen et al., Training Deep Nets with Sublinear Memory Cost
- Inference
- Using a trained model to produce output. There is no backward pass and no optimizer, which is why serving a model costs a fraction of what training it does. The complete inference path
- Interconnect
- The link the GPUs use to talk to each other. It is the hidden limit on every multi-GPU strategy, and it is decided by the machine rather than by the cards.
- Layer
- One processing stage inside the model, taking a list of numbers in and handing a transformed list out. The example model is 28 transformer layers deep. Inside one transformer block
- LoRA
- Low-rank adaptation. Freeze the model and learn a small pair of skinny matrices beside each targeted weight matrix, so about 1 percent of parameters train and the static state for a 1.5B run drops from about 24.6 GB to about 4 GB. Hu et al., LoRA
- Loss
- One number saying how wrong the model was on this batch. Training is the whole business of making it smaller, and a falling loss on its own proves very little. How a neural network learns
- Memory tenant
- One of the four things sharing the card during training: weights, gradients, optimizer state and activations. All four are live at the same moment, so peak memory is their sum rather than the largest of them. The four tenants and the 16 bytes per parameter
- Micro-batch
- The batch that actually fits in memory in one go. Several micro-batches are combined by gradient accumulation so the update behaves like one much larger batch. Batch size, accumulation and the effective batch
- Mixed precision
- Doing the arithmetic in a 16-bit format for speed and memory while keeping a 32-bit copy of the weights so small updates are not lost to rounding. Essentially every modern training run works this way. Micikevicius et al., Mixed Precision Training
- NCCL
- NVIDIA’s library for the collective operations multi-GPU training relies on, such as all-reduce. A handful of its environment variables can silently push traffic onto a slower path with no error message. NVIDIA NCCL user guide
- NVLink
- NVIDIA’s dedicated high-speed link between GPUs, hundreds of gigabytes per second. Recent consumer cards do not have it, so their cross-card traffic goes over PCIe instead.
- Optimizer
- The part of training that turns gradients into actual weight changes. The gradient says which way to move, and the optimizer decides how far. Gradients and optimizers explained
- Optimizer state
- The numbers an optimizer keeps between steps, such as running averages of past gradients. Under standard mixed-precision AdamW it is 12 of the 16 bytes per parameter, which makes it the largest memory tenant. Rajbhandari et al., ZeRO
- Paged optimizer
- An optimizer that can move its state out to ordinary system RAM when GPU memory spikes, then bring it back when the pressure passes. It turns a hard crash into a brief slowdown. Dettmers et al., QLoRA
- Parameter
- One of the numbers inside the model that training can change. Parameter and weight mean the same thing here, and a 1.5B model has about 1.5 billion of them. Training memory and the 16 bytes per parameter
- Parameter-efficient fine-tuning
- Any method that trains a small number of new parameters and leaves the pretrained ones frozen. LoRA and QLoRA are the ones in practical use, and PEFT is the library that implements them. Hugging Face PEFT, LoRA guide
- PCIe
- The general-purpose bus that connects cards to the rest of the machine. Without NVLink it is also how two GPUs talk to each other, at roughly a tenth of the speed and routed through the CPU.
- Pipeline parallelism
- Giving each GPU a different group of layers and passing activations along the chain. It moves far less data than tensor parallelism, at the cost of cards sitting idle while the pipeline fills and drains. Huang et al., GPipe
- Preference tuning
- Training on comparisons rather than single right answers, so a model can be taught that one fluent answer is better than another. It is the stage that installs tone, helpfulness, verbosity control and refusal behaviour. Preference tuning: RLHF, DPO and verifiable rewards
- QLoRA
- LoRA with the frozen model stored in 4 bits instead of 16. It cuts the last remaining cost of simply holding the base weights by about four times, at the price of roughly 40 percent less throughput. Dettmers et al., QLoRA
- Sharding
- Splitting the training state so each GPU holds only a slice of it and fetches the rest when it needs it. It is the fix for a model that will not fit, and it costs traffic between the cards. Rajbhandari et al., ZeRO
- Tensor
- A grid of numbers with any number of dimensions. One number is a scalar, a row of them is a vector, a table is a matrix, and anything past that is still a tensor with more dimensions. How a neural network learns
- Tensor parallelism
- Splitting individual matrix multiplies across GPUs, so partial results have to be combined inside every layer on every pass. It needs a fast link, which makes it a serving technique rather than a training one on consumer cards. Shoeybi et al., Megatron-LM
- Token
- The unit a language model actually reads and writes: a short piece of text, often a word or part of a word. Every length and cost in training is counted in tokens. The complete inference path
- Transformer
- The architecture behind every model in this series: a stack of blocks that alternate attention with a feed-forward network. Vaswani et al., Attention Is All You Need
- Transformer block
- One repeated unit of the model: attention, then a feed-forward network, with normalisation and residual connections around them. A 28-layer model is 28 of these stacked up. Inside one transformer block
- Weight
- A single learned number inside the model, used to multiply an input on its way through a layer. Weights are what fine-tuning changes, and the only thing it changes. What actually changes inside the model
- ZeRO
- The staged sharding scheme: stage one splits the optimizer state, stage two adds the gradients, stage three adds the weights themselves. Use the lowest stage that makes the run fit, because each step up costs more traffic. Rajbhandari et al., ZeRO
Practical exercises
Compute how much slower a gradient sync is on split PCIe than on NVLink
Qwen2.5-1.5B’s gradient in bf16 is 1.54 billion parameters times 2 bytes, about 3.08 GB. Compute how long an all-reduce moving that much data takes over a 900 GB/s NVLink link and over a 64 GB/s link, which is what two consumer cards get once they split a 16-lane PCIe 5.0 slot to 8 lanes each. State the ratio between the two times, and explain why that ratio would come out the same no matter which model’s gradient size you had used.
See the worked solution (opens in a new tab)
Quantify how gradient accumulation cuts interconnect traffic
On a two-card PCIe setup, an all-reduce fires every micro-batch by default, each one moving Qwen2.5-7B’s roughly 14 GB bf16 gradient. If you raise gradient_accumulation_steps from 1 to 8, keeping the micro-batch size fixed, work out how the number of all-reduces per epoch changes and how the total gradient bytes moved per epoch changes. State plainly why the total compute per epoch does not change even though the interconnect traffic does.
See the worked solution (opens in a new tab)
Predict the effect of switching FULL_SHARD to SHARD_GRAD_OP
Take the FSDP YAML from the article and change fsdp_sharding_strategy: FULL_SHARD to fsdp_sharding_strategy: SHARD_GRAD_OP, leaving everything else the same. Using the article’s own mapping between sharding strategies and ZeRO stages, say which memory tenant stops being sharded, what happens to memory usage as a result, and what happens to the volume of cross-card communication. Then state when you would actually prefer this setting over the default.
See the worked solution (opens in a new tab)
Diagnose a distributed job that inherited the wrong NCCL settings
You copy a training container’s environment file from a colleague who ran distributed jobs on datacentre hardware with InfiniBand. On your own desktop with two consumer cards over PCIe, the FSDP job runs, but it is far slower than an identically specced friend’s two-card machine achieves on the same job. Using only the NCCL environment variables described in this part, name the check you would run first, the variable most likely to be the actual cause here, and the one variable in the inherited file that is a red herring rather than the culprit.
See the worked solution (opens in a new tab)
Extend the interconnect reasoning to a four-card workstation
Imagine a workstation with four consumer GPUs sharing one 16-lane PCIe 5.0 root complex, so each card gets roughly 4 lanes instead of the 8 lanes two cards would each get. Using the reasoning from this part, rather than any number specific to four cards, work out, at least as an order-of-magnitude estimate, how per-card bandwidth changes relative to the two-card case. Then give a concrete recommendation for two situations: a job that already fits in a quarter of one card’s memory, and a job that genuinely needs the combined memory of all four cards.
Frequently asked questions
Does using two GPUs let me fine-tune a bigger model?
Only if you shard. Plain data-parallel training puts a full copy of the model and all its training state on every card, so a model that did not fit on one still does not fit on two. Sharded training does raise the ceiling, by splitting the optimizer state, the gradients and the weights across cards, at the cost of gathering them over the interconnect. Even then it divides rather than shrinks: two 32 GB cards do not add up to enough for a 7B full fine-tune.
Do I need NVLink for multi-GPU fine-tuning?
Not for everything. Data-parallel training and sharded training both work over PCIe, they just pay more for it, and gradient accumulation plus adapter-based training cut that cost a long way. What genuinely needs a fast link is tensor-parallel training, which makes the cards talk inside every layer on every pass. Recent consumer cards have no NVLink at all.
What is the difference between DDP and FSDP?
DDP keeps a complete copy of the model and its training state on every card and only synchronises gradients. FSDP splits the weights, gradients and optimizer state across cards, then gathers each layer’s full weights just before computing on it and frees them straight after. The first buys throughput. The second buys memory, for roughly half as much traffic again.
Which ZeRO stage should I use?
The lowest one that makes your run fit. Stage 1 shards the optimizer state, which is 12 of the 16 bytes per parameter, so it delivers the largest saving for almost no extra traffic. Stage 2 adds the gradients for about the same cost. Stage 3 adds the weights and has to gather them layer by layer on every pass, so reach for it only when the first two still leave you short.
Why does tensor parallelism work for inference but not training on the same cards?
Inference is a single forward pass, so the per-layer joining happens once and in one direction, and serving engines are built to overlap and minimise that traffic. Training adds the same pattern in the backward direction on every step, plus gradient and optimizer synchronisation on top. That doubles the most expensive traffic pattern there is, exactly where the interconnect is weakest.
What is the single most effective setting on a slow interconnect?
Gradient accumulation. Raising it means you synchronise once per large effective batch instead of once per micro-batch, which spreads the interconnect cost over far more work for no change in the arithmetic. After that, using an adapter so that only about 1 percent of the gradient bytes cross the wire is the next biggest win.
Sources and further reading
- Rajbhandari et al., ZeRO: Memory Optimizations Toward Training Trillion Parameter Models, the source of both the sharding stages and the 16 bytes per parameter accounting.
- Zhao et al., PyTorch FSDP: Experiences on Scaling Fully Sharded Data Parallel.
- PyTorch, Getting Started with Fully Sharded Data Parallel, now written around the newer FSDP2 API.
- Hugging Face Accelerate, FSDP guide, which documents the config keys above and the mapping from each sharding strategy to its ZeRO stage.
- Shoeybi et al., Megatron-LM, for tensor and pipeline parallelism and how they combine at scale.
- Huang et al., GPipe, on pipeline parallelism and where the bubbles come from.
- Goyal et al., Accurate, Large Minibatch SGD, the linear scaling rule and the gradual warmup that make a large effective batch trainable.
- You et al., Large Batch Optimization for Deep Learning, the LAMB optimizer and the point at which plain linear scaling stops working.
- The NCCL user guide, for the collectives and the environment variables worth checking before a serious run.
