Distribuované trénování v noteboocích

Důležité

Tato funkce je v beta verzi. Správci pracovního prostoru můžou řídit přístup k této funkci ze stránky Previews . Viz Manage Azure Databricks preview.

Dekorátor @distributed ze Serverless GPU Python API je nejpohodlnější způsob, jak spustit distribuovaný trénink v notebooku Databricks. Opatřete svou trénovací funkci dekorátorem, zavolejte ji a AI Runtime ji spustí na všech GPU na uzlu, ke kterému je váš notebook připojen. Stejný kód lze škálovat od jednoho GPU až po více GPU, aniž je nutné zřizovat cluster nebo konfigurovat distribuovaný spouštěč.

Tip

  • Dekorátor @distributed spustí trénovací funkci na všech GPU ve vašem uzlu přímo z notebooku.
  • Podporuje PyTorch DDP, FSDP a DeepSpeed a přesouvá kód s jednou GPU na multi-GPU s minimálními změnami.
  • Připojte notebook k akcelerátoru 8xH100 a nastavte gpus=8 si plný trénink na více GPU.

Rychlý start

Balíček serverless_gpu je předinstalovaný, když je notebook připojen k serverless GPU. Ozdobte svou tréninkovou funkci , @distributeda pak ji .distributed()nazvěte:

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

Každé volání .distributed() vytvoří run MLflow (nebo vnořený podřízený run, pokud už je jeden aktivní) a vypíše odkaz na run ve výstupu buňky. Pro kompletní, průchodný návod viz Plný příklad.

Podporované architektury

Rozhraní @distributed API se integruje s hlavními distribuovanými trénovacími knihovnami:

  • PyTorch Distributed Data Parallel (DDP): Standardní více-GPU datový paralelismus.
  • Plně shardovaný datový paralelismus (FSDP): Paměťově efektivní trénování velkých modelů.
  • DeepSpeed: optimalizační knihovna společnosti Microsoft pro trénování velkých modelů.

Pro skutečné scénáře trénování využívající jednotlivé knihovny viz ukázky notebooků.

Jak dekoratér @distributed funguje

Když zavoláte dekorovanou funkci pomocí .distributed(), AI Runtime spravuje mechaniky, které byste jinak konfigurovali ručně pomocí distribuovaného launcheru:

  • Serializace a rozvětvování: Funkce je serializována a spuštěna na každou z vašich žádostí gpus . Každá GPU spustí kopii funkce se stejnými argumenty.
  • Synchronizace prostředí: Prostředí Python a závislosti jsou replikovány napříč všemi úrovněmi, takže každý proces spouští stejný kód.
  • Řada proměnných prostředí: Standardní proměnné, jako jsou , LOCAL_RANK jsou vyplňovány pro každý proces. Přečtěte si je ve své funkci, abyste umístili model a data na správné zařízení.
  • Sběr výsledků: Hodnoty návratu jsou shromažďovány ze všech pořadí a vráceny volajícímu.
  • Sledování v MLflow: Každé volání .distributed() vytvoří běh v MLflow, nebo vnořený podřízený běh, pokud je již nějaký běh aktivní, takže metriky zaznamenané vaší funkcí se uloží do stejného běhu.
  • Životní cyklus a časový limit: Distribuované provádění probíhá v rámci životního cyklu notebooku. Ukončení poznámkového bloku ukončí běh. Dekoratér má výchozí časovou pauzu 3 hodiny. Zadejte timeout v sekundách, chcete-li to změnit, nebo timeout=None, chcete-li to vypnout. Vlastní časové limity vyžadují GPU prostředí verze 5 a vyšší.

Rozhraní API vychází ze standardních knihoven PyTorch: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) a DeepSpeed.

Přichází z TorchDistributor

Pokud dnes spustíte distribuovaný PyTorch na Sparku s TorchDistributorem a vaše pracovní zátěž se vejde na jeden uzel, serverless_gpu@distributed API je doporučenou náhradou za nové deep learning workloady. Odstraní cluster Sparku a umožní používat stejnou kódovou cestu od jednoho GPU až po více GPU.

funkce serverless_gpu @distributed Rozhraní api Distributor pochodní
Infrastruktura Plně bezserverová, bez správy clusteru Vyžaduje cluster Spark s pracovními procesy GPU.
Setup Jeden dekorátor, minimální konfigurace Vyžaduje cluster Spark a nastavení TorchDistributor.
Podpora rámce PyTorch DDP, FSDP, DeepSpeed Primárně PyTorch DDP
Načítání dat Uvnitř dekorátoru se používají objemy Unity Catalog (UCVolumeDataset pro streamování souborových dat). Přes Spark nebo systém souborů

Chcete-li migrovat jednouchlovou zátěž:

  • Nahraďte volání @distributed dekorátorem train_fn na TorchDistributor(...).run(train_fn, ...) a potom spusťte pomocí train_fn.distributed(...).
  • Odstraňte konfiguraci clusteru Spark a GPU worker. Připojte notebook k akcelerátoru 8xH100 a místo toho nastavte gpus=8.
  • Přesuňte načítání dat dovnitř dekorované funkce. Viz Načítání dat.
  • Uchovejte si stávající kód modelu DDP, FSDP nebo DeepSpeed. Dekoratér podporuje všechny tři.

@distributed běží na jednom uzlu (viz Omezení), takže nenahrazuje všechny pracovní zátěže TorchDistributor. Pracovní zátěže, které závisí na integraci se Sparkem, si nechte na TorchDistributor. Pro spuštění distribuovaného tréninku z vašeho lokálního stroje nebo přes více uzlů použijte místo toho AI Runtime CLI, které je ve veřejném náhledu. Viz AI Runtime CLI.

Úplný příklad

Následující příklad trénuje model vícevrstvého perceptronu (MLP) na 8 GPU H100 v notebooku.

  1. Nastavte model a definujte funkce nástroje.

    
    # 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. Importujte knihovnu serverless_gpudistributed a modul.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Vložte kód pro trénování modelu do funkce a ozdobte funkci dekorátorem @distributed. Dekorovaná funkce je vstupním bodem pro distribuované vykonání, proto v ní definujte veškerou trénovací logiku, načítání dat a inicializaci modelu.

    @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. Spusť distribuované trénování voláním distribuované funkce s uživatelsky definovanými argumenty.

    run_train.distributed(num_epochs=3, batch_size=1)
    
  5. Při spuštění se ve výstupu buňky notebooku vygeneruje odkaz na běh MLflow. Pokud chcete zobrazit výsledky spuštění, klikněte na odkaz spuštění MLflow nebo ho najděte na panelu Experiment . Podrobnosti o přizpůsobení názvů experimentů, sledování metrik a obnovení spuštění najdete v tématu Sledování experimentů a pozorovatelnost.

Načítání dat

Vložte kód pro načítání dat přímo do @distributed funkce. Datová sada může překročit maximální velikost povolenou funkcí pickle, takže její generování nebo načítání v dekorátoru zabrání chybám serializace:

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

Pro data založená na souborech uložená ve volumech Unity Catalog použijte UCVolumeDataset z serverless_gpu.data, který streamuje soubory s lokálním ukládáním do mezipaměti a automaticky je rozděluje mezi ranky a pracovní procesy. Pro uložení kontrolního bodu distribuovaného trénování do svazku použijte UCVolumeWriter a UCVolumeReader. Viz Načítání dat v prostředí AI Runtime a ukládání kontrolních bodů modelu.

Limitations

  • Distribuované trénování běží napříč GPU na jediném uzlu, ke kterému je váš notebook připojen. Pro kompletní trénink s více GPU se připojte k akcelerátoru 8xH100, který poskytuje jednomu uzlu 8 GPU, a nastavte gpus=8.
  • Typ akcelerátoru musí odpovídat. Pokud v @distributed nastavíte gpu_type, musí odpovídat akcelerátoru, ke kterému je váš notebook připojen ("H100" nebo "A10"). Nesoulad způsobuje selhání pracovní zátěže. Tento parametr je volitelný a automaticky detekuje, pokud je vynechán.
  • AI Runtime doporučuje GPU Environment v4 a vyšší. Vlastní časové limity (parametr timeout) vyžadují prostředí GPU verze 5 a vyšší.
  • Dekorátor má ve výchozím nastavení časový limit 3 hodiny. Zadejte timeout v sekundách, chcete-li to změnit, nebo timeout=None, chcete-li to vypnout.
  • Provádění probíhá v rámci životního cyklu notebooku. Ukončení poznámkového bloku ukončí běh.

Další informace