Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
Importante
Este recurso está em versão Beta.
O @distributed decorador da API Serverless GPU Python é a forma mais conveniente de executar treino distribuído a partir de um caderno Databricks. Decora a tua função de treino, chama-a, e o AI Runtime executa-a em todas as GPUs do nó a que o teu portátil está ligado. O mesmo código escala de uma única GPU para várias GPUs, sem necessidade de aprovisionar um cluster nem de configurar um inicializador distribuído.
Pontos-chave nesta página:
- O
@distributeddecorador executa uma função de treino em todas as GPUs do teu nó a partir de dentro de um caderno. - Suporta PyTorch DDP, FSDP e DeepSpeed, e move código de GPU única para multi-GPU com alterações mínimas.
- Ligue o seu portátil a um acelerador 8xH100 e defina
gpus=8para treino multi-GPU completo.
Note
Esta página aborda o treino distribuído em notebooks do Databricks com a API Python serverless para GPU. Para submeter cargas de trabalho de treino distribuídas a partir da sua máquina local, utilize os comandos CLI do Databricks para AI Runtime, que estão em Pré-visualização Pública. Veja Utilizar a CLI do Databricks com runtime de IA.
Início Rápido
O serverless_gpu pacote é pré-instalado quando o seu portátil está ligado a uma GPU serverless. Decore o seu evento de treino com @distributed, e depois chame-o 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 o acelerador H100 8x, que provisiona um único nó com 8 GPUs. Ao usar o @distributed decorador, defina gpus=8. O gpu_type parâmetro é opcional e detetado automaticamente pelo acelerador ao qual o seu portátil está ligado.
Cada chamada a .distributed() cria uma execução do MLflow (ou uma execução secundária aninhada, se já existir uma execução ativa) e imprime uma ligação para a execução na saída da célula. Para um guia completo e executável, veja o exemplo completo.
Estruturas suportadas
A @distributed API integra-se com as principais bibliotecas de treino distribuídas:
- PyTorch Distributed Data Parallel (DDP): Paralelismo padrão de dados multi-GPU.
- Fully Sharded Data Parallel (FSDP): Treino eficiente em memória para modelos de grande dimensão.
- DeepSpeed: A biblioteca de otimização da Microsoft para treino de grandes modelos.
Para ver cenários reais de treino que utilizam cada biblioteca, consulte os exemplos de notebooks.
Como o @distributed decorador trabalha
Ao chamar uma função decorada com .distributed(), o AI Runtime trata dos aspetos técnicos que, de outra forma, configuraria manualmente com um iniciador distribuído:
-
Serialização e distribuição: A função é serializada e executada em cada um dos
gpussolicitados. Cada GPU executa uma cópia da função com os mesmos argumentos. - Sincronização do ambiente: O ambiente Python e as dependências são replicados em todos os níveis, pelo que todos os processos executam o mesmo código.
-
Variáveis de ambiente de classificação: Variáveis padrão, como as
LOCAL_RANKque são preenchidas para cada processo. Leia-as na sua função para colocar o modelo e os dados no dispositivo correto. - Recolha de resultados: Os valores de retorno são recolhidos de todas as classificações e devolvidos ao chamador.
-
Rastreio do MLflow: Cada chamada
.distributed()cria uma execução do MLflow, ou uma execução filha aninhada, se já houver uma ativa, pelo que as métricas registadas pela sua função são associadas à mesma execução. -
Ciclo de vida e tempo limite: A execução distribuída ocorre durante o ciclo de vida do notebook. Ao terminar o notebook, a execução é terminada. O decorador tem um tempo limite predefinido de 3 horas. Introduza
timeoutem segundos para o alterar, outimeout=Nonepara o desativar. Os tempos limite personalizados requerem os ambientes GPU v5 ou superior.
A API baseia-se nas bibliotecas padrão do PyTorch: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) e DeepSpeed.
Vindo do TorchDistributor
Se hoje executas o PyTorch distribuído no Spark com o TorchDistributor e a tua carga de trabalho cabe num único nó, a serverless_gpu@distributed API é a alternativa recomendada para novas cargas de deep learning. Remove o cluster Spark e proporciona o mesmo percurso de código de uma única GPU para várias GPUs.
| Feature |
serverless_gpu
@distributed API |
TorchDistribuidor |
|---|---|---|
| Infraestrutura | Totalmente serverless, sem gestão de clusters | Requer um cluster Spark com trabalhadores da GPU |
| Configuração | Decorador único, configuração mínima | Requer configuração do cluster Spark e do TorchDistributor. |
| Suporte à estrutura | PyTorch DDP, FSDP, DeepSpeed | Principalmente PyTorch DDP |
| Carregamento de dados | No decorador, utiliza volumes do Unity Catalog (UCVolumeDataset para transmissão contínua de dados de ficheiros) |
Via Spark ou sistema de ficheiros |
Para migrar uma carga de trabalho de nó único:
- Substitua a chamada
TorchDistributor(...).run(train_fn, ...)pelo decorador@distributedemtrain_fne, em seguida, inicie comtrain_fn.distributed(...). - Remova o cluster Spark e a configuração do trabalhador da GPU. Liga o teu portátil a um acelerador 8xH100 e configura
gpus=8em vez disso. - Mover o carregamento de dados dentro da função decorada. Ver carregamento de dados.
- Mantém o código existente do modelo DDP, FSDP ou DeepSpeed. O decorador apoia os três.
@distributed corre num único nó (ver Limitações), por isso não substitui todas as cargas de trabalho do TorchDistributor. Mantém as cargas de trabalho que dependem da integração com o Spark no TorchDistributor. Para executar formação distribuída a partir da sua máquina local ou em vários nós, utilize, em alternativa, a CLI do Runtime de IA, que está em Pré-visualização pública. Veja Utilizar a CLI do Databricks com runtime de IA.
Exemplo completo
O exemplo seguinte treina um modelo de perceptrão multicamada (MLP) em 8 GPUs H100 a partir de um portátil.
Configura o teu modelo e define funções de utilidade.
# 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)Importa a
serverless_gpubiblioteca e odistributedmódulo.import serverless_gpu from serverless_gpu import distributedEnvolve o código de treino do modelo numa função e decora a função com o
@distributeddecorador. A função decorada é o ponto de entrada para a execução distribuída, por isso define toda a lógica de treino, carregamento de dados e inicialização do modelo 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 treino distribuído chamando a função distribuída com argumentos definidos pelo utilizador.
run_train.distributed(num_epochs=3, batch_size=1)Quando é executado, é gerada uma ligação para uma execução no MLflow na saída da célula do notebook. Clique no link da execução do MLflow ou encontre-o no painel Experimentos para ver os resultados da execução. Para detalhes sobre como personalizar nomes de experiências, rastrear métricas e retomar execuções, consulte Rastreio e observabilidade de experiências.
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, pelo que 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 ficheiros armazenados em volumes do Unity Catalog, use UCVolumeDataset de databricks.air.data, que transmite ficheiros com cache local e os distribui automaticamente entre ranks e workers. Para criar um ponto de verificação do treino distribuído num volume, use UCVolumeWriter e UCVolumeReader. Veja Carregar dados no AI Runtime e Checkpoint com Distributed Checkpoint (DCP).
Limitations
- O treino distribuído é executado pelas GPUs do nó único ao qual o seu notebook está ligado. Para treino multi-GPU completo, ligue-se a um acelerador 8xH100, que disponibiliza um nó com 8 GPUs, e defina
gpus=8. - O tipo de acelerador deve corresponder. Se definires
gpu_typeem@distributed, deve corresponder ao acelerador a que o teu portátil está ligado ("H100"ou"A10"). Uma incompatibilidade faz falhar a carga de trabalho. O parâmetro é opcional e deteta-se automaticamente quando é omitido. - O AI Runtime recomenda o ambiente GPU v4 e superiores. Os limites de tempo personalizados (parâmetro
timeout) requerem o ambiente GPU v5 ou superior. - Por predefinição, o decorador expira ao fim de 3 horas. Introduza
timeoutem segundos para o alterar, outimeout=Nonepara o desativar. - A execução ocorre no ciclo de vida do notebook. Ao terminar o notebook, a execução é terminada.
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 o seu pipeline de treino mais eficiente e resiliente, consulte o guia de desempenho e resiliência.
- Para cenários de treino de ponta a ponta, veja exemplos de cadernos.