Formação distribuída em cadernos

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 @distributed decorador 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=8 para 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 gpus solicitados. 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_RANK que 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 timeout em segundos para o alterar, ou timeout=None para 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 @distributed em train_fn e, em seguida, inicie com train_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=8 em 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.

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

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Envolve o código de treino do modelo numa função e decora a função com o @distributed decorador. 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()
    
  4. 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)
    
  5. 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_type em @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 timeout em segundos para o alterar, ou timeout=None para o desativar.
  • A execução ocorre no ciclo de vida do notebook. Ao terminar o notebook, a execução é terminada.

Recursos adicionais