Training distribuito nei notebook

Importante

Questa funzionalità è in versione beta. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.

Il @distributed decoratore dell'API Python Serverless GPU è il modo più pratico per eseguire l'addestramento distribuito da un notebook di Databricks. Decora la tua funzione di addestramento, chiamala, e AI Runtime la esegue su tutte le GPU del nodo a cui è collegato il tuo notebook. Lo stesso codice scala da una sola GPU a multi-GPU senza cluster da fornire e senza launcher distribuito da configurare.

Tip

  • Il decoratore @distributed esegue una funzione di addestramento su tutte le GPU del nodo all’interno di un notebook.
  • Supporta PyTorch DDP, FSDP e DeepSpeed, e sposta il codice da singola GPU su multi-GPU con modifiche minime.
  • Collega il tuo notebook a un acceleratore 8xH100 e imposta gpus=8 per un addestramento multi-GPU completo.

Quickstart

Il serverless_gpu pacchetto viene preinstallato quando il notebook è collegato a una GPU serverless. Decora la tua attività di allenamento con @distributed, poi chiamala con .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()

Ogni chiamata a .distributed() crea una run di MLflow (o una run figlia annidata, se è già attiva una run) e stampa un collegamento alla run nell'output della cella. Per una guida completa ed eseguibile, vedi l'esempio completo.

Framework supportati

L'API @distributed si integra con le principali librerie di training distribuite:

  • PyTorch Distributed Data Parallel (DDP): Parallelismo dei dati standard su più GPU.
  • Fully Sharded Data Parallel (FSDP): addestramento efficiente in termini di memoria per modelli di grandi dimensioni.
  • DeepSpeed: libreria di ottimizzazione di Microsoft per il training di modelli di grandi dimensioni.

Per scenari di addestramento reali che utilizzano ogni libreria, vedi gli esempi dei quaderni.

Come lavora il @distributed decoratore

Quando chiami una funzione decorata con .distributed(), AI Runtime gestisce le meccaniche che altrimenti configureresti manualmente con un launcher distribuito:

  • Serializzazione e dispersione: La funzione viene serializzata e avviata su ciascuna delle gpus richieste che richiedi. Ogni GPU esegue una copia della funzione con gli stessi argomenti.
  • Sincronizzazione dell'ambiente: L'ambiente Python e le dipendenze sono replicati su tutti i livelli, quindi ogni processo esegue lo stesso codice.
  • Variabili di ambiente di rango: variabili standard come LOCAL_RANK quelle vengono popolate per ogni processo. Leggili nella tua funzione per posizionare il modello e i dati sul dispositivo corretto.
  • Raccolta dei risultati: I valori restituiti vengono raccolti da tutti i ranghi e restituiti al chiamante.
  • Tracciamento MLflow: Ogni .distributed() chiamata crea una run MLflow, o una run figlia annidata se già attiva, quindi le metriche registrate dalla tua funzione finiscono sulla stessa run.
  • Ciclo di vita e timeout: L'esecuzione distribuita si svolge all'interno del ciclo di vita del notebook. Terminare il quaderno interrompe la run. Il decoratore ha un timeout predefinito di 3 ore. Passa timeout in pochi secondi per cambiarlo, o timeout=None per disattivarlo. I timeout personalizzati richiedono l'ambiente GPU v5 e superiore.

L'API si basa sulle librerie standard di PyTorch: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) e DeepSpeed.

Proveniente da TorchDistributor

Se oggi esegui PyTorch distribuito su Spark con TorchDistributor e il tuo carico di lavoro si adatta a un singolo nodo, l'API serverless_gpu@distributed è il sostituto raccomandato per i nuovi carichi di lavoro di deep learning. Rimuove il cluster Spark e ti dà lo stesso percorso di codice da una sola GPU a multi-GPU.

Feature serverless_gpu @distributed API TorchDistributor
Infrastruttura Completamente serverless, nessuna gestione del cluster Richiede un cluster Spark con ruoli di lavoro GPU
Setup Decorator singolo, configurazione minima Richiede l'installazione di cluster Spark e TorchDistributor
Supporto per framework PyTorch DDP, FSDP, DeepSpeed Principalmente PyTorch DDP
Caricamento dei dati Nel decoratore, utilizza i volumi di Unity Catalog (UCVolumeDataset per lo streaming dei dati di file) Tramite Spark o filesystem

Per migrare un carico di lavoro a singolo nodo:

  • Sostituisci la chiamata a TorchDistributor(...).run(train_fn, ...) con il decoratore @distributed in train_fn, quindi avvia con train_fn.distributed(...).
  • Rimuovi il cluster Spark e la configurazione del worker GPU. Collega il tuo notebook a un acceleratore 8xH100 e imposta gpus=8 invece.
  • Sposta il caricamento dei dati all'interno della funzione decorata. Vedi caricamento dati.
  • Conserva il codice dei modelli DDP, FSDP o DeepSpeed esistente. Il decoratore sostiene tutti e tre.

@distributed gira su un singolo nodo (vedi Limitazioni), quindi non sostituisce ogni carico di lavoro TorchDistributor. Mantieni i carichi di lavoro che dipendono dall'integrazione con Spark su TorchDistributor. Per eseguire l'addestramento distribuito dalla tua macchina locale o su più nodi, usa invece la CLI AI Runtime, che si trova in Public Preview. Vedi AI Runtime CLI.

Esempio completo

Il seguente esempio addestra un modello multilayer perceptron (MLP) su 8 GPU H100 da un notebook.

  1. Configurare il modello e definire le funzioni di utilità.

    
    # 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. Importare la serverless_gpu libreria e il distributed modulo.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Includere il codice di addestramento del modello in una funzione e decorare la funzione con il decoratore @distributed. La funzione decorata è il punto di ingresso per l'esecuzione distribuita, quindi definisci tutta la logica di addestramento, il caricamento dei dati e l'inizializzazione dei modelli al suo interno.

    @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. Esegui l'addestramento distribuito chiamando la funzione distribuita con argomenti definiti dall'utente.

    run_train.distributed(num_epochs=3, batch_size=1)
    
  5. Quando viene eseguito, nell'output della cella del notebook viene generato un collegamento di esecuzione MLflow. Fare clic sul collegamento MLflow run o trovarlo nel pannello Experiment per visualizzare i risultati dell'esecuzione. Per informazioni dettagliate sulla personalizzazione dei nomi degli esperimenti, il rilevamento delle metriche e la ripresa delle esecuzioni, vedere Rilevamento e osservabilità dell'esperimento.

Caricamento dati

Inserisci il codice di caricamento dati all'interno della @distributed funzione. Un dataset può superare la dimensione massima consentita da pickle, quindi generandolo o caricandolo all'interno del decorator si evitano errori di serializzazione:

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

Per i dati basati su file archiviati nei volumi di Unity Catalog, usare UCVolumeDataset da serverless_gpu.data, che trasmette i file con memorizzazione nella cache locale e li partiziona automaticamente tra i ranghi e i ruoli di lavoro. Per eseguire il checkpoint di training distribuito a un volume, usare UCVolumeWriter e UCVolumeReader. Vedere Caricare i dati nel runtime di intelligenza artificiale e nel checkpoint del modello.

Limitations

  • L'addestramento distribuito si svolge tra le GPU sul singolo nodo a cui è collegato il tuo notebook. Per un addestramento completo multi-GPU, collegati a un acceleratore 8xH100, che fornisce un nodo con 8 GPU, e imposta gpus=8.
  • Il tipo di acceleratore deve corrispondere. Se imposta gpu_type , @distributeddeve corrispondere all'acceleratore a cui il tuo taccuino è collegato ("H100" oppure "A10"). Una mancata corrispondenza provoca il fallimento del carico di lavoro. Il parametro è opzionale e viene rilevato automaticamente quando omesso.
  • AI Runtime consiglia l'ambiente GPU v4 o versioni successive. I timeout personalizzati (il timeout parametro) richiedono l'ambiente GPU v5 e superiore.
  • Per impostazione predefinita, il decoratore va in timeout dopo 3 ore. Passa timeout in pochi secondi per cambiarlo, o timeout=None per disattivarlo.
  • L'esecuzione si svolge all'interno del ciclo di vita del notebook. La terminazione del notebook interrompe l'esecuzione.

Ulteriori informazioni