Escalar LLMs com paralelismo multi-GPU e multi-nó
Tradução automática Este artigo foi traduzido automaticamente a partir da versão original em inglês.
As cargas de trabalho com modelos de grande dimensão ultrapassam uma única GPU por diferentes motivos. Um job de treino pode ficar sem memória devido ao estado do optimizer. Outro pode esgotar a memória com as activations de sequências longas. Um modelo que cabe na memória pode, ainda assim, não atingir o throughput pretendido. Cada problema exige uma partição e um padrão de comunicação diferentes.
Este é um guia prático das principais estratégias de paralelismo e das restrições que lhes estão subjacentes, baseado no Ultra-Scale Playbook da Hugging Face. O objectivo é mostrar o que cada divisão permite ganhar, que dados comunica e quando se tornam necessárias combinações.
Em resumo. O paralelismo de dados replicado aumenta o throughput de treino quando uma réplica cabe na memória. O fully sharded data parallelism particiona o estado do modelo, mas acrescenta all-gathers de parâmetros e reduce-scatters de gradients. O tensor, o pipeline, o contexto e o expert parallelism dividem, respectivamente, a matemática das layers, a profundidade, a sequência e as layers de mixture-of-experts (MoE). Combine-os apenas depois de identificar a restrição determinante de memória ou comunicação.
Este guia pressupõe que está familiarizado com backpropagation, layers Transformer e um loop de treino PyTorch standard.
Comece por dois orçamentos de memória
O treino e a inferência não têm o mesmo footprint.
training peak ≈ parameters
+ gradients
+ optimizer state
+ saved activations
+ temporary buffers
+ communication buffers
+ allocator headroom
inference peak ≈ resident weights
+ key/value (KV) cache
+ runtime workspace
+ communication buffers
+ allocator headroom
Um modelo com 70 mil milhões de parâmetros tem um limite inferior decimal de 140 GB apenas para os pesos BF16. Este valor diz pouco sobre o treino, onde os gradients, o estado do optimizer, os master weights e as activations podem dominar. Também não dimensiona o serving, onde a política de cache, o comprimento da sequência, a concorrência de batches e a quantização são relevantes.
Faça profiling da arquitectura exacta, da precisão, do comprimento da sequência, do micro-batch, do optimizer, da política de checkpointing e do runtime. Registe a memória máxima alocada e reservada, os tokens por segundo, o tempo passado em kernels e o tempo exposto em collectives.
Cada dimensão de paralelismo implica um compromisso
Para cada estratégia, pergunte qual a dimensão do tensor que é dividida, que estado é replicado e que collective entra no critical path.
| Estratégia | Divide | Principal benefício | Comunicação introduzida |
|---|---|---|---|
| Paralelismo de dados replicado | batch | throughput de treino | all-reduce de gradients |
| Fully sharded data parallelism | parâmetros, gradients e estado do optimizer num grupo de data parallel (DP) | memória do estado do modelo | all-gather de parâmetros, reduce-scatter de gradients |
| Tensor parallelism | dimensões de matrizes ou de attention dentro das layers | pesos e activations das layers | collectives dentro dos transformer blocks |
| Pipeline parallelism | grupos de layers | profundidade do modelo e estado por stage | activations ponto a ponto e bubbles de scheduling |
| Context parallelism | dimensão da sequência | memória de activations de sequências longas | troca de key/value ou de attention entre o grupo da sequência |
| Expert parallelism | experts MoE e tokens encaminhados | capacidade de experts por rank | dispatch e combinação de tokens, normalmente all-to-all |
O benefício em memória não é um multiplicador fixo. Depende do grau de sharding, do que permanece replicado, do estado temporariamente não sharded, da política de activations, do padding, do desequilíbrio e dos buffers.
Paralelismo de dados replicado: throughput sem capacidade adicional
O paralelismo de dados replicado, normalmente designado por distributed data parallel, mantém uma réplica completa do treino em cada rank. A documentação DistributedDataParallel do PyTorch descreve este modelo de replicação e sincronização de gradients. Cada rank processa um micro-batch diferente e os gradients são sincronizados antes do optimizer step.
Use-o quando todo o estado do treino couber com margem de segurança e o batch global puder aumentar, ou quando for possível ajustar a gradient accumulation. As principais vantagens são semântica simples e uma implementação madura que sobrepõe o cálculo backward à redução de gradients em buckets.
Adicionar ranks pode ser prejudicial quando o batch local se torna demasiado pequeno. Também pode prejudicar o desempenho quando a rede não consegue ocultar o all-reduce ou quando há stalls na entrega dos dados de entrada. O batch de optimização pretendido pode não escalar.
Fully sharded data parallelism: memória de estado em troca de collectives
O fully sharded data parallelism armazena shards dos parâmetros, gradients e optimizer num grupo. O artigo sobre ZeRO descreve este padrão de sharding do estado e a API FSDP2 do PyTorch implementa-o. Os parâmetros de uma layer são reunidos através de all-gather para o cálculo e podem voltar a ser sharded posteriormente. Os gradients são devolvidos aos owners através de reduce-scatter.
A documentação actual do PyTorch distingue a API fully_shard, no fully sharded data parallelism versão 2 (FSDP2), do wrapper FullyShardedDataParallel mais antigo. O FSDP2 agrupa a comunicação pelos módulos aos quais fully_shard é aplicado e recomenda uma aplicação bottom-up, para que os grupos de layers possam sobrepor comunicação e cálculo.
from torch.distributed.fsdp import fully_shard
from torch.optim import AdamW
# Apply bottom-up: each block becomes a communication group.
for block in model.transformer.blocks:
fully_shard(block)
# Shard remaining root parameters such as embeddings and output projection.
fully_shard(model)
# Construct the optimizer after parameters have become sharded distributed tensors (DTensors).
optimizer = AdamW(model.parameters(), lr=learning_rate)
Este é um esquema estrutural, não um launcher completo. Device meshes, mixed precision, checkpointing, inicialização, estado do optimizer e distributed checkpoints têm de ser compatíveis com a stack de treino.
O sharding é atractivo quando o estado do modelo é a restrição determinante e o cálculo das layers consegue ocultar tráfego collective suficiente. Pode ser uma má troca para modelos pequenos, links lentos, layers pequenas ou layouts cujo grupo de shards atravessa a fronteira de topologia errada.
Tensor parallelism: particionar a matemática das layers
O tensor parallelism particiona a álgebra linear dentro de uma layer. São exemplos as projections column-parallel e row-parallel. O guia de paralelismo da NVIDIA documenta esta divisão ao nível das layers. Os resultados parciais exigem collectives dentro dos transformer blocks, pelo que a latência e a largura de banda são relevantes repetidamente durante os passes forward e backward.
Use-o quando uma layer ou as suas activations não couberem, ou quando as matmuls forem suficientemente grandes para que kernels particionados superem um único rank. Mapeie o grupo tensor-parallel para o domínio de comunicação mais rápido disponível e depois faça medições. Um grau elevado de tensor parallelism pode reduzir cada matriz local até ao ponto em que a eficiência dos kernels diminui, enquanto o overhead dos collectives aumenta.
O sequence parallelism é frequentemente combinado com tensor parallelism para evitar a replicação de algum trabalho sobre as activations. É distinto do context parallelism, que opera sobre a sequência de entrada completa do modelo.
Pipeline parallelism: particionar a profundidade e o tempo do scheduling
O pipeline parallelism coloca grupos de layers diferentes em stages distintos e envia activations entre eles. Os micro-batches mantêm os stages a trabalhar em simultâneo. O artigo sobre GPipe usa este scheduling para redes neuronais gigantes.
Reduz o estado do modelo por stage e pode diminuir o volume de comunicação que atravessa uma fronteira mais lenta, em comparação com collectives tensor-parallel por layer. Os custos incluem bubbles, transferências de activations, desequilíbrio entre stages, scheduling mais complexo e recuperação e checkpointing mais difíceis.
Num schedule ao estilo GPipe simples e equilibrado, com p stages e m micro-batches, a fracção idealizada do bubble forward é aproximadamente:
(p - 1) / (m + p - 1)
Os schedules reais podem usar variantes one-forward/one-backward, interleaving ou zero-bubble, e o custo desigual das layers pode dominar a fórmula. Escolha os limites dos stages a partir de tempos e memória medidos, não de uma divisão igual do número de layers.
Context parallelism: particionar as activations de sequências longas
O context parallelism distribui a dimensão da sequência. A documentação de context parallelism da NVIDIA descreve a divisão da sequência e a troca de key/value necessária para a attention. Cada rank é responsável por um shard da sequência, enquanto a attention troca a informação necessária para preservar a semântica de contexto completo. As implementações podem usar rings ponto a ponto, all-gather, all-to-all ou combinações hierárquicas.
Reduz a memória de activations no treino com contexto longo, mas replica os pesos no grupo de contexto e introduz comunicação de attention. O benefício depende do tipo de attention, da máscara causal, do comprimento da sequência, do recomputation e da forma como os grupos de contexto se combinam com os grupos tensor e data parallel.
Não o seleccione com base num limiar universal de 8K, 32K ou 100K. Faça profiling da memória de activations e da comunicação de attention para a arquitectura real.
Expert parallelism: apenas para uma arquitectura MoE
O expert parallelism distribui os experts nas layers mixture-of-experts. O guia de paralelismo da NVIDIA documenta esta colocação de experts e a sua combinação com outras dimensões de paralelismo. O router envia as representações dos tokens para os experts seleccionados e combina os resultados. Apenas os experts seleccionados calculam para cada token, mas o total de pesos dos experts continua a exigir armazenamento e colocação adequada no serving.
O expert parallelism não é um interruptor de optimização para um modelo dense. Faz parte de uma arquitectura MoE. As preocupações incluem equilíbrio de carga, limites de capacidade, all-to-all de tokens, tokens descartados ou preenchidos com padding, auxiliary losses e desequilíbrio de falhas. Monitorize tokens por expert, entropia do routing, overflow de capacidade, tempo de comunicação e qualidade por route.
Compor um layout a partir da topologia
Os sistemas de treino de modelos dense de grande dimensão sem um grupo expert-parallel usam normalmente um produto dos tamanhos dos grupos data-parallel (DP), tensor-parallel (TP), pipeline-parallel (PP) e context-parallel (CP):
world size = DP × TP × PP × CP
Quando o expert parallelism (EP) é um grupo independente, o guia de paralelismo da NVIDIA calcula o total como:
total GPUs = TP × PP × CP × EP × DP
Use a mesh suportada pelo framework, em vez de multiplicar configurações não suportadas.
Construa o layout pela seguinte ordem:
- Desenhe os domínios de comunicação: links GPU-a-GPU, switches, fronteiras de non-uniform memory access (NUMA), fabric do nó, oversubscription e caminho de storage.
- Coloque os collectives frequentes e sensíveis à latência, normalmente TP, no domínio adequado mais rápido.
- Escolha grupos FSDP ou DP replicados a partir da capacidade e da largura de banda restantes.
- Adicione PP quando a colocação por profundidade ou o tráfego entre domínios trouxer benefícios, equilibrando o tempo medido e a memória dos stages.
- Adicione CP apenas para a restrição da sequência e expert parallelism (EP) apenas para a topologia de experts do modelo.
- Confirme a divisibilidade de heads, dimensões hidden, layers, experts, batch e sequência para a mesh candidata.
- Faça benchmark de várias meshes válidas. As heurísticas conscientes da topologia escolhem candidatos, não vencedores.
Dois clusters com o mesmo número de GPUs podem preferir layouts diferentes, porque a largura de banda dos links, a hierarquia de switches, a ligação dos CPUs e a contenção da rede são diferentes.
O treino e o serving exigem decisões separadas
A inferência normalmente não transporta gradients nem estado do optimizer, pelo que os layouts de treino ao estilo FSDP não são automaticamente transferíveis.
Para o serving, pergunte:
- Uma réplica consegue conter os pesos, o KV cache, o workspace e a concorrência pretendida?
- O throughput é melhor servido por mais réplicas independentes ou por sharding de uma única réplica?
- O TP reduz suficientemente a pressão sobre os pesos e o cache por rank para compensar a comunicação por layer?
- O PP é suportado de forma eficiente pelo modelo e pelo request scheduler?
- Como é que o prefill e o decode pressionam de forma diferente o cálculo, a largura de banda da memória e o interconnect?
- O que acontece à tail latency quando os requests têm comprimentos diferentes de prompt e output?
Faça benchmark do servidor completo com o scheduler, a quantização, a distribuição do contexto, a política de batching e o perfil de tráfego. Os tokens por segundo do treino não permitem prever o time to first token nem a latência entre tokens no serving.
Medir honestamente um layout de scaling
Para cada candidato, registe:
- identidade do modelo, código, runtime, kernels e topologia
- batch global e local, distribuição das sequências e contagem de tokens
- memória máxima por categoria, quando disponível
- tokens úteis por segundo e utilização de operações de vírgula flutuante do modelo (FLOP), quando calculados de forma consistente
- tempo exposto em operações all-reduce, all-gather, reduce-scatter, all-to-all e ponto a ponto
- stalls de entrada, tempo de checkpoint, comportamento após restart e distribuição dos stragglers
- loss de treino ou paridade dos outputs do serving face ao baseline
Compare deliberadamente strong scaling e weak scaling. O strong scaling mantém o trabalho total constante à medida que aumenta o número de ranks. O weak scaling mantém constante o trabalho por rank, pelo que o trabalho total cresce com o número de ranks. Uma percentagem identificada como “eficiência de scaling” não tem significado sem esse denominador e baseline.
Conclusão
O paralelismo é um mapeamento entre um bottleneck medido, uma dimensão do tensor e um padrão de comunicação. A replicação, o sharding, a partição de layers, o staging, a partição da sequência e o routing de experts aliviam restrições diferentes e criam modos de falha diferentes.
Faça o inventário da carga de trabalho, desenhe a topologia, gere meshes válidas e faça profiling. O layout vencedor é aquele que cabe com margem e minimiza a comunicação exposta para o job que realmente executa.
Referências
- Documentação do PyTorch FSDP2
fully_shard- Sharding de parâmetros ao nível dos módulos, all-gathers, reduce-scatters e comportamento de DTensor. - Documentação DistributedDataParallel do PyTorch - Treino de modelos replicados com gradients sincronizados.
- Estratégias de paralelismo do NVIDIA Megatron Core - Dimensões de data, tensor, pipeline, contexto, experts e fully sharded parallelism.
- Hugging Face Ultra-Scale Playbook - Visão geral interactiva do model parallelism em grande escala.
- Artigo sobre DeepSpeed ZeRO - Sharding de estados do optimizer, gradients e parâmetros para reduzir a utilização de memória.
- Artigo sobre GPipe - Treino pipeline com micro-batches e execução em stages.
- Documentação de paralelismo da NERSC - Definições de strong e weak scaling baseadas em trabalho total ou por rank fixo.