Treinamento distribuído em notebooks

Importante

Esse recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

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.

Dica

  • O decorador @distributed executa 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=8 para treinamento completo multi-GPU.

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()

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 espalhamento: A função é serializada e lançada em cada uma das gpus solicitações que você solicita. 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-as na 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 timeout em segundos para mudar ou timeout=None desativar. 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 No decorador, utiliza volumes do Unity Catalog (UCVolumeDataset para streaming de dados de arquivo) Via Spark ou sistema de arquivos

Para migrar uma carga de trabalho de nó único:

  • Substitua a chamada @distributed pelo decorador train_fn em train_fn.distributed(...), depois inicie com TorchDistributor(...).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=8 em 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 o treinamento distribuído na sua máquina local ou em vários nós, use a CLI do AI Runtime, que está em Public Preview. Consulte a CLI do AI Runtime.

Exemplo completo

O exemplo a seguir treina um modelo de perceptron multicamada (MLP) em 8 GPUs H100 a partir de um notebook.

  1. 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)
    
  2. Importe a serverless_gpu biblioteca e o distributed módulo.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Encapsule 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()
    
  4. 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)
    
  5. 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 serverless_gpu.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. Consulte Carregar dados em Runtime de IA e Ponto de verificação do modelo.

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_type em @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 timeout em segundos para mudar ou timeout=None desativar.
  • A execução é executada dentro do ciclo de vida do notebook. Encerrar o notebook encerra a execução.

Saiba mais