Распределённое обучение в ноутбуках

Это важно

Эта функция доступна в бета-версии. Администраторы рабочей области могут управлять доступом к этой функции на странице "Предварительные версии ". См. статью "Управление предварительными версиями Azure Databricks".

Декоратор @distributed из Serverless GPU Python API — это самый удобный способ запустить распределённое обучение из ноутбука Databricks. Оформите функцию обучения, вызовите её, и AI Runtime запускает её на всех GPU на узле, к которому подключён ваш ноутбук. Один и тот же код масштабируется от одного GPU до нескольких GPU без необходимости разворачивать кластер и настраивать средство распределённого запуска.

Tip

  • @distributed Декоратор запускает обучающую функцию для каждого GPU на вашем узле из блокнота.
  • Он поддерживает PyTorch DDP, FSDP и DeepSpeed, а также переносит код с одной GPU в мульти-GPU с минимальными изменениями.
  • Подключите ноутбук к ускорителю 8xH100 и настройте gpus=8 для полноценного обучения на нескольких GPU.

Note

Эта страница посвящена распределенному обучению в ноутбуках Databricks с использованием бессерверного Python API для GPU. Чтобы отправить распределённые обучающие нагрузки с вашего локального компьютера, используйте команды Databricks CLI для AI Runtime, которые находятся в публичном предпросмотре. См. Использование CLI Databricks с AI Runtime.

Quickstart

serverless_gpu Пакет предустановлен, когда ноутбук подключён к бессерверной видеокарте. Добавьте декоратор к обучающей функции с помощью @distributed, затем вызовите её с помощью .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()

Для распределенного обучения требуется акселератор 8xH100, который подготавливает один узел с 8 GPU. При использовании декоратора @distributed установите gpus=8. Параметр gpu_type необязателен и автоматически обнаруживается из акселератора, к которому подключена записная книжка.

Каждый вызов .distributed() создаёт запуск MLflow (или вложенный дочерний запуск, если уже есть активный запуск) и выводит ссылку в выводе ячейки. Полный рабочий пример см. в разделе полный пример.

Поддерживаемые платформы

@distributed API интегрируется с основными распределенными библиотеками обучения:

  • PyTorch Distributed Data Parallel (DDP): стандартный параллелизм данных с несколькими GPU.
  • Полностью шардированный параллелизм данных (FSDP): обучение больших моделей с эффективным использованием памяти.
  • DeepSpeed: библиотека оптимизации Microsoft для обучения больших моделей.

Для реальных обучающих сценариев, использующих каждую библиотеку, см. примеры в блокнотах.

Как @distributed работает декоратор

Когда вы вызываете декорированную функцию с помощью .distributed(), AI Runtime берёт на себя всю техническую сторону, которую в противном случае пришлось бы настраивать вручную с помощью средства запуска распределённых задач:

  • Сериализация и распространение: функция сериализируется и запускается по каждому запросу gpus . Каждая видеокарта запускает копию функции с одними и теми же аргументами.
  • Синхронизация среды: Среда Python и зависимости реплицируются на всех рангах, поэтому каждый процесс выполняет один и тот же код.
  • Переменные среды ранга: стандартные переменные, такие как , LOCAL_RANK заполняются для каждого процесса. Прочитайте их в своей функции, чтобы разместить модель и данные на правильном устройстве.
  • Сбор результатов: Возвращаемые значения собираются со всех рангов и возвращаются вызывающему.
  • Отслеживание MLflow: Каждый .distributed() вызов создаёт MLflow run или вложенный дочерний запуск, если он уже активен, поэтому метрики, зарегистрированные из вашей функции, попадают в тот же запуск.
  • Жизненный цикл и тайм-аут: Распределённое выполнение происходит в рамках жизненного цикла ноутбука. При остановке блокнота выполнение также прерывается. У декоратора по умолчанию тайм-аут составляет 3 часа. Укажите timeout в секундах, чтобы изменить это, или timeout=None, чтобы отключить это. Кастомные тайм-ауты требуют GPU среды v5 и выше.

API основан на стандартных библиотеках PyTorch: Distributed Data Parallel (DDP), Fully Sharded Data Parallel (FSDP) и DeepSpeed.

Поступает от TorchDistributor

Если вы сегодня запускаете распределённый PyTorch на Spark с TorchDistributor и ваша нагрузка помещается на одном узле, serverless_gpu@distributed API является рекомендуемой заменой для новых задач глубокого обучения. Это устраняет необходимость в кластере Spark и обеспечивает единый путь выполнения кода от одного GPU до нескольких GPU.

Функция serverless_gpu @distributed API Распространитель факелов
Инфраструктура Полностью бессерверное решение, без управления кластерами Требуется кластер Spark со службами, работающими на GPU.
Setup Один декоратор, минимальная конфигурация Требуется настройка кластера Spark и TorchDistributor
Поддержка фреймворка PyTorch DDP, FSDP, DeepSpeed В первую очередь PyTorch DDP
Загрузка данных Внутри декоратора используются тома Unity Catalog (UCVolumeDataset для потоковой передачи файлов) С помощью Spark или файловой системы

Для миграции одноузловой рабочей нагрузки:

  • Замените вызов TorchDistributor(...).run(train_fn, ...) на декоратор @distributed для train_fn, затем запустите с помощью train_fn.distributed(...).
  • Удалите кластер Spark и конфигурацию GPU worker. Подключите ноутбук к ускорителю 8xH100 и вместо этого задайте gpus=8.
  • Переместите загрузку данных внутрь декорированной функции. См. загрузку данных.
  • Сохраняйте существующий код модели DDP, FSDP или DeepSpeed. Декоратор поддерживает все три варианта.

@distributed работает на одном узле (см. Ограничения), поэтому не заменяет все рабочие нагрузки TorchDistributor. Оставьте рабочие нагрузки, требующие интеграции со Spark, на TorchDistributor. Для запуска распределённого обучения с локального компьютера или между несколькими узлами используйте AI Runtime CLI, который находится в публичном предпросмотре. См. Использование CLI Databricks с AI Runtime.

Полный пример

Следующий пример обучает многослойную модель перцептрона (MLP) на 8 графических процессорах H100 из ноутбука.

  1. Настройте модель и определите служебные функции.

    
    # 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_gpudistributed и модуль.

    import serverless_gpu
    from serverless_gpu import distributed
    
  3. Оберните код обучения модели в функцию и декорируйте функцию декоратором @distributed. Украшенная функция является входной точкой распределённого выполнения, поэтому задайте всю обучающую логику, загрузку данных и инициализацию модели внутри неё.

    @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. Запускайте распределённое обучение, вызывая распределённую функцию с пользовательскими аргументами.

    run_train.distributed(num_epochs=3, batch_size=1)
    
  5. При выполнении в выходных данных ячейки блокнота появляется ссылка на запуск MLflow. Щелкните ссылку запуска MLflow или найдите ее на панели "Эксперимент ", чтобы просмотреть результаты выполнения. Дополнительные сведения о настройке имен экспериментов, отслеживания метрик и возобновлении выполнения см. в разделе "Отслеживание экспериментов" и "Наблюдаемость".

Загрузка данных

Поместите код загрузки данных внутри функции @distributed. Набор данных может превышать максимальный размер, допустимый для pickle, поэтому создание или загрузка его внутри декоратора позволяет избежать ошибок сериализации:

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

Для файловых данных, хранящихся в томах Unity Catalog, используйте UCVolumeDataset из serverless_gpu.data, который передаёт файлы в потоковом режиме с локальным кэшированием и автоматически распределяет их между рангами и рабочими процессами. Для контрольных точек распределенного обучения в том, используйте UCVolumeWriter и UCVolumeReader. См. Загрузка данных в AI Runtime и Checkpoint with Distributed Checkpoint (DCP).

Limitations

  • Распределённое обучение проходит по GPU на одном узле, к которому подключён ваш ноутбук. Для полного обучения с использованием нескольких GPU подключитесь к ускорителю 8xH100, который выделяет один узел с 8 GPU, и установите gpus=8.
  • Тип ускорителя должен совпадать. Если вы задаёте gpu_type в @distributed, оно должно соответствовать акселератору, к которому подключён ваш ноутбук ("H100" или "A10"). Несоответствие приводит к сбою рабочей нагрузки. Параметр необязательный и автоматически определяется при его отсутствии.
  • AI Runtime рекомендует использовать среду с GPU v4 и выше. Настраиваемые тайм-ауты (параметр timeout) требуют среды GPU версии v5 или выше.
  • По умолчанию тайм-аут декоратора составляет 3 часа. Укажите timeout в секундах, чтобы изменить это, или timeout=None, чтобы отключить это.
  • Выполнение происходит в рамках жизненного цикла ноутбука. При остановке блокнота выполнение также прерывается.

Узнать больше