Distribuerad träning i anteckningsböcker

Viktigt!

Den här funktionen finns i Beta. Arbetsyteadministratörer kan styra åtkomsten till den här funktionen från sidan Förhandsversioner . Se Hantera förhandsversioner av Azure Databricks.

@distributed-dekoratorn från Serverless GPU Python API är det smidigaste sättet att köra distribuerad träning från en Databricks-notebook. Dekorera din träningsfunktion, anropa den, och AI Runtime kör den på alla GPU:er på den nod som din notebook är ansluten till. Samma kod skalar från en GPU till flera GPU:er utan kluster att provisionera och utan distribuerad launcher att konfigurera.

Tip

  • Dekoratören @distributed kör en träningsfunktion på alla GPU:er på din nod från en anteckningsboksmiljö.
  • Den stöder PyTorch DDP, FSDP och DeepSpeed, och flyttar kod från ett enda GPU till flera GPU med minimala ändringar.
  • Anslut din notebook till en 8xH100-accelerator och ställ in gpus=8 för full multi-GPU-träning.

Snabbstart

Paketet serverless_gpu är förinstallerat när din laptop är ansluten till ett serverlöst grafikkort. Dekorera din träningsfunktion med @distributed, och kalla den sedan med .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()

Varje anrop av .distributed() skapar en MLflow-körning (eller en kapslad underordnad körning om det redan finns en aktiv körning) och visar en länk till körningen i cellens utdata. För en komplett, körbar genomgång, se Fullständigt exempel.

Ramverk som stöds

API:et @distributed integreras med större distribuerade utbildningsbibliotek:

  • PyTorch Distributed Data Parallel (DDP): Standard multi-GPU data parallellism.
  • Fullständigt fragmenterad dataparallell (FSDP): Minneseffektiv träning för stora modeller.
  • DeepSpeed: Microsoft optimeringsbibliotek för stor modellträning.

För verkliga träningsscenarier som använder varje bibliotek, se exempel på anteckningsböcker.

Hur dekoratören @distributed arbetar

När du anropar en dekorerad funktion med .distributed(), hanterar AI Runtime det praktiska som du annars skulle behöva konfigurera manuellt med ett distribuerat startprogram:

  • Serialisering och förgrening: Funktionen serialiseras och startas på var och en av de gpus som du begär. Varje GPU kör en kopia av funktionen med samma argument.
  • Miljösynkronisering: Python-miljön och beroendena replikeras över alla nivåer, så varje process kör samma kod.
  • Miljövariabler för rank: Standardvariabler som LOCAL_RANK fylls i för varje process. Läs dem i din funktion för att placera modellen och data på rätt enhet.
  • Resultatinsamling: Returvärden samlas in från alla nivåer och återlämnas till den som ringer.
  • MLflow-spårning: Varje .distributed() anrop skapar en MLflow-körning, eller en nästlad barnkörning om en redan är aktiv, så att mätvärden som loggas från din funktion hamnar på samma körning.
  • Livscykel och timeout: Distribuerad exekvering körs inom anteckningsbokens livscykel. Om du avslutar anteckningsboken avslutas körningen. Dekoratorn har en standardtidsgräns på 3 timmar. Passera timeout på några sekunder för att ändra det, eller timeout=None inaktivera det. Anpassade tidsgränser kräver GPU-miljö v5 eller senare.

API:et bygger vidare på de standardiserade PyTorch-biblioteken: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) och DeepSpeed.

Kommer från TorchDistributor

Om du kör distribuerad PyTorch på Spark idag med TorchDistributor och din arbetsbelastning passar på en enda nod, är API:et serverless_gpu@distributed den rekommenderade ersättningen för nya djupinlärningsarbetsbelastningar. Den tar bort Spark-klustret och ger dig samma kodväg från enskilt grafikkort till fler-GPU.

Feature serverless_gpu @distributed API Lyktdistributör
Infrastruktur Fullständigt serverlös, ingen klusterhantering Kräver ett Spark-kluster med GPU-arbetare
Inställningar Enkel dekoratör, minimal konfiguration Kräver konfiguration av Spark-kluster och TorchDistributor
Stöd för ramverk PyTorch DDP, FSDP, DeepSpeed Främst PyTorch DDP
Datainläsning Inuti dekoratorn används Unity-katalogvolymer (UCVolumeDataset för strömmande fildata). Via Spark eller filsystem

För att migrera en arbetsbelastning med en enda nod:

  • Byt ut anropet TorchDistributor(...).run(train_fn, ...) mot dekoratören @distributedtrain_fn, och starta sedan med train_fn.distributed(...).
  • Ta bort Spark-klustret och GPU-worker-konfigurationen. Koppla din notebook till en 8xH100-accelerator och ställ in gpus=8 istället.
  • Flytta data som laddas in i den dekorerade funktionen. Se Data loading.
  • Behåll din befintliga DDP-, FSDP- eller DeepSpeed-modellkod. Dekoratören har stöd för alla tre.

@distributed körs på en enda nod (se Begränsningar), så den ersätter inte varje TorchDistributor-arbetsbelastning. Behåll arbetsbelastningar som är beroende av Spark-integration på TorchDistributor. För att köra distribuerad träning från din lokala maskin eller över flera noder, använd istället AI Runtime CLI, som finns i Public Preview. Se AI Runtime CLI.

Fullständigt exempel

Följande exempel tränar en flerlagerperceptronmodell (MLP) på 8 H100-GPU:er i en notebook.

  1. Konfigurera din modell och definiera verktygsfunktioner.

    
    # 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. serverless_gpu Importera biblioteket och modulendistributed.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Lägg in modellträningskoden i en funktion och använd dekoratören @distributed på funktionen. Den dekorerade funktionen är ingångspunkten för distribuerad exekvering, så definiera all träningslogik, dataladdning och modellinitiering inuti den.

    @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. Kör den distribuerade träningen genom att anropa den distribuerade funktionen med användardefinierade argument.

    run_train.distributed(num_epochs=3, batch_size=1)
    
  5. När den körs genereras en MLflow-körningslänk i notebook-cellens utdata. Klicka på MLflow-körningslänken eller hitta den i panelen Experiment för att se körningsresultatet. Mer information om hur du anpassar experimentnamn, spårar mått och återupptar körningar finns i Experimentspårning och observerbarhet.

Dataladdning

Placera dataladdningskod i @distributed funktionen. En datamängd kan överstiga den maximala storlek som tillåts av pickle, så att generera eller ladda in den i dekoratören undviker serialiseringsfel:

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

För filbaserade data som lagras i Unity Catalog-volymer använder du UCVolumeDataset från serverless_gpu.data, som strömmar filer med lokal cachelagring och partitionerar dem mellan led och arbetare automatiskt. Om du vill spara kontrollpunkter för distribuerad träning till en volym använder du UCVolumeWriter och UCVolumeReader. Se Läsa in data om AI-körning och modellkontrollpunkter.

Limitations

  • Distribuerad träning körs över GPU:erna på den enda nod som din notebook är ansluten till. För full multi-GPU-träning, koppla till en 8xH100-accelerator, som tillhandahåller en nod med 8 GPU:er, och sätt gpus=8.
  • Typen av accelerator måste överensstämma. Om du anger gpu_type i @distributed måste det stämma överens med den accelerator som din notebook är ansluten till ("H100" eller "A10"). En konflikt leder till att arbetslasten misslyckas. Parametern är valfri och upptäcks automatiskt när den utelämnas.
  • AI Runtime rekommenderar GPU-miljö v4 eller senare. Anpassade tidsgränser (parametern timeout) kräver GPU-miljö v5 eller senare.
  • Dekoratören går ut efter 3 timmar som standard. Passera timeout på några sekunder för att ändra det, eller timeout=None inaktivera det.
  • Exekveringen sker inom anteckningsbokens livscykel. Om du avslutar anteckningsboken avslutas körningen.

Lära sig mer