Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Importante
Esse recurso está em Beta.
O decorador @distributed da API Serverless GPU para Python é a forma mais conveniente de executar treinamento distribuído em um notebook do Databricks. Decore sua função de treinamento, chame-a, e o AI Runtime roda em todas as GPUs do nó ao qual seu notebook está conectado. O mesmo código escala de GPU única para multiGPU, sem cluster para provisionar e sem launcher distribuído para configurar.
Pontos-chave desta página:
- O decorador
@distributedexecuta uma função de treinamento em cada GPU do seu nó de dentro de um notebook. - Ele suporta PyTorch DDP, FSDP e DeepSpeed, e move código de GPU única para multi-GPU com mudanças mínimas.
- Conecte seu notebook a um acelerador 8xH100 e configure
gpus=8para treinamento completo multi-GPU.
Note
Esta página aborda o treinamento distribuído a partir de notebooks do Databricks com a API Python de GPU sem servidor. Para enviar cargas de trabalho de treinamento distribuídas da sua máquina local, use os comandos da CLI do Databricks para AI Runtime, que estão no Public Preview. Veja Usar a CLI do Databricks com o AI Runtime.
Início Rápido
O serverless_gpu pacote é pré-instalado quando seu notebook está conectado a uma GPU serverless. Decore sua função de treinamento com @distributed e depois chame-a com .distributed():
from serverless_gpu import distributed
# gpus is the number of GPUs on the node. gpu_type is optional and
# auto-detected from the accelerator your notebook is connected to.
@distributed(gpus=8, gpu_type="H100")
def train():
import os
import torch
import torch.distributed as dist
# Bind this process to its own GPU before training.
local_rank = int(os.environ["LOCAL_RANK"])
torch.cuda.set_device(local_rank)
device = torch.device(f"cuda:{local_rank}")
dist.init_process_group("nccl")
# ... build the model and data on `device`, then run your training loop ...
dist.destroy_process_group()
train.distributed()
O treinamento distribuído requer um acelerador 8xH100, que provisiona um único nó com 8 GPUs. Ao usar o @distributed decorador, defina gpus=8. O gpu_type parâmetro é opcional e detectado automaticamente no acelerador ao qual seu notebook está conectado.
Cada chamada de .distributed() cria uma execução no MLflow (ou uma execução filha aninhada, se já houver uma execução ativa) e exibe um link da execução na saída da célula. Para um guia completo e executável, veja o exemplo completo.
Estruturas com suporte
A @distributed API se integra às principais bibliotecas de treinamento distribuídas:
- DDP (Distributed Data Parallel) do PyTorch: paralelismo de dados de várias GPUs padrão.
- FSDP (Fully Sharded Data Parallel): treinamento com eficiência de memória para modelos grandes.
- DeepSpeed: biblioteca de otimização do Microsoft para treinamento de modelo grande.
Para cenários reais de treinamento que usam cada biblioteca, veja exemplos de cadernos.
Como o @distributed decorador trabalha
Quando você chama uma função decorada com .distributed(), o AI Runtime lida com os mecanismos que, de outra forma, você configuraria manualmente com um inicializador distribuído:
-
Serialização e fan-out: a função é serializada e executada em cada um dos
gpusque você solicitar. Cada GPU executa uma cópia da função com os mesmos argumentos. - Sincronização de ambientes: O ambiente Python e as dependências são replicados em todos os níveis, então todo processo executa o mesmo código.
-
Variáveis de ambiente de rank: Variáveis padrão, como
LOCAL_RANK, são definidas para cada processo. Leia-os em sua função para colocar o modelo e os dados no dispositivo correto. - Coleta de resultados: Os valores de retorno são coletados de todas as classificações e devolvidos ao chamador.
-
Rastreamento do MLflow: Cada chamada
.distributed()cria uma execução do MLflow, ou uma execução filha aninhada, se já houver uma execução ativa, de modo que as métricas registradas por sua função sejam associadas à mesma execução. -
Ciclo de vida e timeout: A execução distribuída roda dentro do ciclo de vida do notebook. Encerrar o notebook encerra a execução. O decorador tem um timeout padrão de 3 horas. Passe
timeoutem segundos para alterá-lo outimeout=Nonepara desativá-lo. Limites de tempo personalizados exigem o ambiente GPU v5 ou superior.
A API se baseia nas bibliotecas padrão do PyTorch: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) e DeepSpeed.
Vindo da TorchDistributor
Se você executa o PyTorch distribuído no Spark hoje com TorchDistributor e sua carga de trabalho cabe em um único nó, a API serverless_gpu@distributed é a substituição recomendada para novas cargas de trabalho de aprendizado profundo. Ele remove o cluster Spark e fornece o mesmo fluxo de código de uma única GPU para várias GPUs.
| Característica |
serverless_gpu
@distributed API |
Distribuidor de Tochas |
|---|---|---|
| Infraestrutura | Totalmente sem servidor, sem gerenciamento de cluster | Requer um cluster Spark com trabalhadores de GPU |
| Configuração | Decorador único, configuração mínima | Requer a configuração do cluster Spark e do TorchDistributor |
| Suporte para estrutura | PyTorch DDP, FSDP, DeepSpeed | Principalmente PyTorch DDP |
| Carregamento de dados | Dentro do decorador, utiliza volumes do Catálogo do Unity (UCVolumeDataset para dados de arquivos de streaming) |
Via Spark ou sistema de arquivos |
Para migrar uma carga de trabalho de nó único:
- Substitua a chamada
@distributedpelo decoradortrain_fnemtrain_fn.distributed(...), depois inicie comTorchDistributor(...).run(train_fn, ...). - Remova o cluster do Spark e a configuração do trabalhador da GPU. Conecte seu notebook a um acelerador 8xH100 e defina
gpus=8em vez disso. - Mova o carregamento de dados dentro da função decorada. Veja Carregamento de dados.
- Mantenha seu código de modelo DDP, FSDP ou DeepSpeed existente. O decorador apoia os três.
@distributed é executado em um único nó (consulte Limitações), portanto não substitui todas as cargas de trabalho do TorchDistributor. Mantenha cargas de trabalho que dependem da integração com o Spark no TorchDistributor. Para executar treinamento distribuído a partir de sua máquina local ou em vários nós, use a CLI do AI Runtime, que está em Versão preliminar pública. Veja Usar a CLI do Databricks com o AI Runtime.
Exemplo completo
O exemplo a seguir treina um modelo de perceptron multicamada (MLP) em 8 GPUs H100 a partir de um notebook.
Configure seu modelo e defina funções utilitárias.
# Define the model import os import torch import torch.distributed as dist import torch.nn as nn def setup(): torch.cuda.set_device(int(os.environ["LOCAL_RANK"])) dist.init_process_group("nccl") def cleanup(): dist.destroy_process_group() class SimpleMLP(nn.Module): def __init__(self, input_dim=10, hidden_dim=64, output_dim=1): super().__init__() self.net = nn.Sequential( nn.Linear(input_dim, hidden_dim), nn.ReLU(), nn.Dropout(0.2), nn.Linear(hidden_dim, hidden_dim), nn.ReLU(), nn.Dropout(0.2), nn.Linear(hidden_dim, output_dim) ) def forward(self, x): return self.net(x)Importe a
serverless_gpubiblioteca e odistributedmódulo.import serverless_gpu from serverless_gpu import distributedEncapsule o código de treinamento do modelo em uma função e decore a função com o decorador
@distributed. A função decorada é o ponto de entrada para a execução distribuída, então defina toda a lógica de treinamento, carregamento de dados e inicialização de modelos dentro dela.@distributed(gpus=8, gpu_type='H100') def run_train(num_epochs: int, batch_size: int) -> None: import mlflow import torch.optim as optim from torch.nn.parallel import DistributedDataParallel as DDP from torch.utils.data import DataLoader, DistributedSampler, TensorDataset # 1. Set up multi-GPU environment setup() device = torch.device(f"cuda:{int(os.environ['LOCAL_RANK'])}") # 2. Apply the Torch distributed data parallel (DDP) library for data-parellel training. model = SimpleMLP().to(device) model = DDP(model, device_ids=[device]) # 3. Create and load dataset. x = torch.randn(5000, 10) y = torch.randn(5000, 1) dataset = TensorDataset(x, y) sampler = DistributedSampler(dataset) dataloader = DataLoader(dataset, sampler=sampler, batch_size=batch_size) # 4. Define the training loop. optimizer = optim.Adam(model.parameters(), lr=0.001) loss_fn = nn.MSELoss() for epoch in range(num_epochs): sampler.set_epoch(epoch) model.train() total_loss = 0.0 for step, (xb, yb) in enumerate(dataloader): xb, yb = xb.to(device), yb.to(device) optimizer.zero_grad() loss = loss_fn(model(xb), yb) # Log loss to MLflow metric mlflow.log_metric("loss", loss.item(), step=step) loss.backward() optimizer.step() total_loss += loss.item() * xb.size(0) mlflow.log_metric("total_loss", total_loss) print(f"Total loss for epoch {epoch}: {total_loss}") cleanup()Execute o treinamento distribuído chamando a função distribuída com argumentos definidos pelo usuário.
run_train.distributed(num_epochs=3, batch_size=1)Quando executado, um link de execução do MLflow é gerado na saída da célula de notebook. Clique no link de execução do MLflow ou localize-o no painel Experimento para ver os resultados da execução. Para obter detalhes sobre como personalizar nomes de experimentos, controlar métricas e retomar execuções, consulte Acompanhamento e observabilidade de experimentos.
Carregamento de dados
Coloque o código de carregamento de dados dentro da @distributed função. Um conjunto de dados pode exceder o tamanho máximo permitido por pickle, então gerá-lo ou carregá-lo dentro do decorador evita erros de serialização:
from serverless_gpu import distributed
# This may cause a pickle error because the dataset is captured by the function.
dataset = get_dataset(file_path)
@distributed(gpus=8, gpu_type='H100')
def run_train():
# Load the dataset inside the decorated function instead.
dataset = get_dataset(file_path)
...
Para dados baseados em arquivo armazenados em volumes do Unity Catalog, use UCVolumeDataset de databricks.air.data, que transmite arquivos com armazenamento em cache local e os particiona automaticamente entre classificações e trabalhos. Para fazer o ponto de verificação do treinamento distribuído em um volume, use UCVolumeWriter e UCVolumeReader. Veja Carregar dados no AI Runtime e Checkpoint com Checkpoint Distribuído (DCP).
Limitations
- O treinamento distribuído é executado nas GPUs do único nó ao qual seu notebook está conectado. Para treinamento completo com múltiplas GPUs, conecte-se a um acelerador 8xH100, que fornece um nó com 8 GPUs, e configure
gpus=8. - O tipo de acelerador deve corresponder. Se você definir
gpu_typeem@distributed, ele deve corresponder ao acelerador ao qual o notebook está conectado ("H100"ou"A10"). Uma incompatibilidade faz a carga de trabalho falhar. O parâmetro é opcional e é detectado automaticamente quando omitido. - O AI Runtime recomenda o ambiente GPU v4 e superiores. Timeouts personalizados (o parâmetro
timeout) exigem ambientes de GPU v5 ou superior. - O decorador expira após 3 horas por padrão. Passe
timeoutem segundos para alterá-lo outimeout=Nonepara desativá-lo. - A execução é executada dentro do ciclo de vida do notebook. Encerrar o notebook encerra a execução.
Recursos adicionais
- Para o decorador
@distributed,GPUTypee as APIs do Ray, consulte a documentação de referência da API Python de GPU sem servidor. - Para padrões que tornam seu pipeline de treinamento mais eficiente e resiliente, veja o guia de desempenho e resiliência.
- Para cenários de treinamento de ponta a ponta, veja exemplos de cadernos.