Pelatihan terdistribusi dalam buku catatan

Penting

Fitur ini ada di Beta. Admin ruang kerja dapat mengontrol akses ke fitur ini dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.

@distributed dekorator dari Serverless GPU Python API adalah cara yang paling praktis untuk menjalankan pelatihan terdistribusi melalui notebook Databricks. Tambahkan dekorator pada fungsi pelatihan Anda, lalu panggil fungsi tersebut, dan AI Runtime akan menjalankannya pada semua GPU di node tempat notebook Anda terhubung. Kode yang sama dapat diskalakan dari single-GPU ke multi-GPU tanpa perlu menyediakan cluster dan tanpa perlu mengonfigurasi launcher terdistribusi.

Tip

  • Dekorator @distributed menjalankan fungsi pelatihan pada setiap GPU pada node Anda dari dalam notebook.
  • Mendukung PyTorch DDP, FSDP, dan DeepSpeed, serta memindahkan kode GPU tunggal ke multi-GPU dengan perubahan minimal.
  • Hubungkan notebook Anda ke akselerator 8xH100 dan atur gpus=8 untuk pelatihan multi-GPU penuh.

Quickstart

serverless_gpu Paket sudah terpasang sebelumnya ketika notebook Anda terhubung ke GPU serverless. Dekorasi fungsi latihan Anda dengan @distributed, lalu sebut dengan .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()

Setiap pemanggilan .distributed() membuat run MLflow (atau run turunan bertingkat jika sudah ada run yang aktif) dan menampilkan tautan run pada output sel. Untuk panduan lengkap yang dapat dijalankan, lihat Contoh lengkap.

Kerangka kerja yang didukung

@distributed API terintegrasi dengan pustaka pelatihan terdistribusi utama:

  • PyTorch Distributed Data Parallel (DDP): Paralelisme data multi-GPU standar.
  • Fully Sharded Data Parallel (FSDP): Pelatihan hemat memori untuk model besar.
  • DeepSpeed: pustaka pengoptimalan Microsoft untuk pelatihan model besar.

Untuk skenario pelatihan nyata yang menggunakan setiap pustaka, lihat contoh buku catatan.

Cara kerja dekorator @distributed

Saat Anda memanggil fungsi yang didekorasi dengan .distributed(), AI Runtime menangani mekanik yang biasanya Anda konfigurasikan secara manual dengan peluncur terdistribusi:

  • Serialisasi dan penyebaran: Fungsi ini diserialkan dan dijalankan pada setiap gpus yang Anda minta. Setiap GPU menjalankan salinan fungsi dengan argumen yang sama.
  • Sinkronisasi lingkungan: Lingkungan Python dan dependensi direplikasi di semua tingkatan, sehingga setiap proses menjalankan kode yang sama.
  • Variabel lingkungan rank: Variabel standar seperti LOCAL_RANK ditetapkan untuk setiap proses. Bacalah semuanya dalam fungsi Anda untuk menempatkan model dan data pada perangkat yang tepat.
  • Pengumpulan hasil: Nilai hasil dikumpulkan dari semua peringkat dan dikembalikan ke penelepon.
  • Pelacakan MLflow: Setiap .distributed() panggilan membuat MLflow run, atau nested child run jika sudah aktif, sehingga metrik yang dicatat dari fungsi Anda akan berada di run yang sama.
  • Siklus hidup dan batas waktu: Eksekusi terdistribusi berjalan selama siklus hidup buku catatan. Mengakhiri notebook akan mengakhiri proses berjalan. Dekorator memiliki waktu istirahat default selama 3 jam. Lewatkan timeout dalam hitungan detik untuk mengubahnya, atau timeout=None menonaktifkannya. Timeout khusus memerlukan lingkungan GPU versi 5 atau yang lebih baru.

API ini dibangun di atas pustaka PyTorch standar: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP), dan DeepSpeed.

Berasal dari TorchDistributor

Jika Anda menjalankan PyTorch terdistribusi di Spark saat ini dengan TorchDistributor dan beban kerja Anda dapat ditangani pada satu node, API serverless_gpu@distributed adalah pengganti yang direkomendasikan untuk beban kerja pembelajaran mendalam yang baru. Perangkat ini menghapus klaster Spark dan memberikan jalur kode yang sama dari single-GPU ke multi-GPU.

Feature serverless_gpu @distributed API TorchDistributor
Infrastruktur Sepenuhnya tanpa server, tidak ada manajemen kluster Memerlukan kluster Spark dengan pekerja GPU
Siapkan Dekorator tunggal, konfigurasi minimal Memerlukan kluster Spark dan pengaturan TorchDistributor
Dukungan kerangka kerja PyTorch DDP, FSDP, DeepSpeed Terutama PyTorch DDP
Pemuatan data Dalam dekorator, menggunakan volume Unity Catalog (UCVolumeDataset untuk streaming data berkas) Melalui Spark atau sistem file

Untuk memigrasikan beban kerja satu node:

  • Ganti pemanggilan TorchDistributor(...).run(train_fn, ...) dengan dekorator @distributed pada train_fn, lalu jalankan dengan train_fn.distributed(...).
  • Hapus konfigurasi cluster Spark dan pekerja GPU. Hubungkan notebook Anda ke akselerator 8xH100 dan atur gpus=8 sebagai gantinya.
  • Pindahkan pemuatan data ke dalam fungsi yang didekorasi. Lihat Pemuatan Data.
  • Pertahankan kode model DDP, FSDP, atau DeepSpeed yang sudah ada. Dekorator mendukung ketiganya.

@distributed berjalan pada satu node (lihat Keterbatasan), sehingga tidak menggantikan setiap beban kerja TorchDistributor. Simpan beban kerja yang bergantung pada integrasi Spark di TorchDistributor. Untuk menjalankan pelatihan terdistribusi dari mesin lokal Anda atau di beberapa node, gunakan AI Runtime CLI sebagai gantinya, yang berada di Public Preview. Lihat AI Runtime CLI.

Contoh lengkap

Contoh berikut melatih model multilayer perceptron (MLP) menggunakan 8 GPU H100 melalui notebook.

  1. Siapkan model Anda dan tentukan fungsi utilitas.

    
    # 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. Impor pustaka serverless_gpu dan modul distributed.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Bungkus kode pelatihan model dalam fungsi dan dekori fungsi dengan dekorator @distributed. Fungsi yang didekorasi adalah titik masuk untuk eksekusi terdistribusi, jadi definisikan semua logika pelatihan, pemuatan data, dan inisialisasi model di dalamnya.

    @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. Jalankan pelatihan terdistribusi dengan memanggil fungsi terdistribusi menggunakan argumen yang ditentukan pengguna.

    run_train.distributed(num_epochs=3, batch_size=1)
    
  5. Saat dijalankan, tautan proses MLflow dibuat pada output sel notebook. Klik tautan eksekusi MLflow atau temukan di panel Eksperimen untuk melihat hasil eksekusi. Untuk detail tentang menyesuaikan nama eksperimen, melacak metrik, dan melanjutkan eksekusi, lihat Pelacakan dan pengamatan eksperimen.

Pemuatan data

Tempatkan kode pemuatan data di dalam @distributed fungsi. Sebuah dataset dapat melebihi ukuran maksimum yang diizinkan oleh pickle, sehingga menghasilkan atau memuatnya di dalam dekorator menghindari kesalahan serialisasi:

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

Untuk data berbasis file yang disimpan dalam volume Unity Catalog, gunakan UCVolumeDataset dari serverless_gpu.data, yang melakukan streaming file dengan cache lokal dan membagi file-file tersebut ke seluruh rank dan pekerja secara otomatis. Untuk memeriksa pelatihan terdistribusi ke volume, gunakan UCVolumeWriter dan UCVolumeReader. Lihat Memuat data pada Runtime AI dan Titik pemeriksaan Model.

Limitations

  • Pelatihan terdistribusi berjalan di seluruh GPU pada satu node yang terhubung dengan notebook Anda. Untuk pelatihan multi-GPU penuh, hubungkan ke akselerator 8xH100, yang menyediakan satu node dengan 8 GPU, dan atur gpus=8.
  • Tipe akselerator harus cocok. Jika Anda menetapkan gpu_type di @distributed, nilainya harus sesuai dengan akselerator yang terhubung ke notebook Anda ("H100" atau "A10"). Ketidakcocokan menyebabkan beban kerja gagal. Parameter ini bersifat opsional dan otomatis terdeteksi jika dihilangkan.
  • AI Runtime merekomendasikan lingkungan GPU v4 ke atas. Batas waktu khusus (parameter timeout) memerlukan lingkungan GPU v5 atau yang lebih baru.
  • Dekorator biasanya habis waktu setelah 3 jam secara default. Lewatkan timeout dalam hitungan detik untuk mengubahnya, atau timeout=None menonaktifkannya.
  • Eksekusi berlangsung dalam siklus hidup notebook. Mengakhiri notebook akan mengakhiri proses berjalan.

Pelajari lebih lanjut