Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
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
@distributedkö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=8fö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
gpussom 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_RANKfylls 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
timeoutpå några sekunder för att ändra det, ellertimeout=Noneinaktivera 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@distributedpåtrain_fn, och starta sedan medtrain_fn.distributed(...). - Ta bort Spark-klustret och GPU-worker-konfigurationen. Koppla din notebook till en 8xH100-accelerator och ställ in
gpus=8istä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.
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)serverless_gpuImportera biblioteket och modulendistributed.import serverless_gpu from serverless_gpu import distributedLägg in modellträningskoden i en funktion och använd dekoratören
@distributedpå 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()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)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_typei@distributedmå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
timeoutpå några sekunder för att ändra det, ellertimeout=Noneinaktivera det. - Exekveringen sker inom anteckningsbokens livscykel. Om du avslutar anteckningsboken avslutas körningen.
Lära sig mer
- För
@distributeddecorator-,GPUType, och Ray-API:erna, se referensdokumentationen Serverless GPU Python API. - För mönster som gör din träningspipeline mer effektiv och motståndskraftig, se guiden för prestation och resiliens.
- För träningsscenarier från början till slut, se notebookexempel.