ノートブックでの分散トレーニング

Important

この機能は ベータ版です。 ワークスペース管理者は、[ プレビュー] ページからこの機能へのアクセスを制御できます。 Manage Azure Databricks プレビューを参照してください。

Serverless GPU Python API@distributedデコレーターは、Databricksノートブックから分散トレーニングを実行する最も便利な方法です。 トレーニング関数を装飾して呼び出すと、AI Runtimeがノートブックに接続されているノード上のすべてのGPUで実行します。 同じコードは、クラスタプロビジョニングや分散ランチャーの設定なしに、シングルGPUからマルチGPUへとスケールします。

ヒント

  • @distributedデコレーターはノートパソコンの中からノード上のすべてのGPUに対してトレーニング機能を実行します。
  • PyTorch DDP、FSDP、DeepSpeedをサポートし、最小限の変更でシングルGPUコードをマルチGPUに移行します。
  • ノートパソコンを 8xH100 アクセラレータに接続し、 gpus=8 をマルチGPUトレーニングに設定してください。

クイック スタート

ノートパソコンがサーバーレスGPUに接続されている場合、 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()

.distributed() を呼び出すたびに、MLflow ランが作成されます(すでにアクティブなランがある場合は、ネストされた子ランが作成されます)。また、セルの出力にランへのリンクが表示されます。 完全で実行可能なウォークスルーについては、 Full例をご覧ください。

サポートされているフレームワーク

@distributed API は、主要な分散トレーニング ライブラリと統合されます。

  • PyTorch 分散データ並列 (DDP): 標準のマルチ GPU データ並列処理。
  • 完全シャーディングデータ並列 (FSDP): 大規模モデル向けのメモリ効率に優れたトレーニング。
  • DeepSpeed: 大規模なモデル トレーニング用のMicrosoftの最適化ライブラリ。

各ライブラリを使った実際のトレーニングシナリオについては、 ノートブックの例を参照してください。

@distributedデコレーターの仕組み

.distributed()で装飾関数を呼び出すと、AI Runtimeは分散ランチャーで手動で設定するような仕組みを処理します。

  • シリアライズとファンアウト:関数はシリアライズされ、リクエストされた各 gpus で起動されます。 すべてのGPUは同じ引数で関数のコピーを実行します。
  • 環境同期:Python環境と依存関係はすべてのランクで複製されているため、すべてのプロセスが同じコードを実行します。
  • 環境変数のランク付け: LOCAL_RANK などの標準変数が各プロセスごとに入力されます。 関数でそれらを読み、モデルとデータを正しいデバイスに配置してください。
  • 結果収集:すべてのランクから返される値は収集され、コールプレイヤーに返されます。
  • MLflowのトラッキング:各 .distributed() 呼び出しはMLflowランを生成し、すでにアクティブな場合はネストされた子ランを作成します。これにより、関数から記録された指標が同じ実行に集まります。
  • ライフサイクルとタイムアウト:分散実行はノートブックのライフサイクル内で実行されます。 ノートブックを終了するとランは終了します。 デコレーターのデフォルトのタイムアウトは3時間です。 数秒で timeout をパスして変更するか、 timeout=None で無効にしてください。 カスタムタイムアウトにはGPU環境v5以上が必要です。

APIは標準的なPyTorchライブラリである Distributed Data Parallel (DDP)、 Fully Sharded Data Parallel (FSDP)、 DeepSpeedを基盤としています。

TorchDistributorより

もし現在Spark上で TorchDistributor を使って分散型PyTorchを実行し、ワークロードが1つのノードに収まるなら、 serverless_gpu@distributed APIは新しいディープラーニングワークロードの推奨代替となります。 Sparkクラスターを外し、シングルGPUからマルチGPUへのコードパスが同じになります。

特徴 serverless_gpu @distributed API TorchDistributor
インフラストラクチャ 完全にサーバーレス、クラスター管理なし GPU ワーカーを使用する Spark クラスターが必要
セットアップ 単一デコレータ、最小構成 Spark クラスターと TorchDistributor のセットアップが必要
フレームワーク サポート PyTorch DDP、FSDP、DeepSpeed 主に PyTorch DDP
データの読み込み デコレーター内部ではUnity Catalogボリューム(ストリーミングファイルデータ用UCVolumeDataset )を使用しています Spark またはファイルシステム経由

単一ノードワークロードの移行:

  • TorchDistributor(...).run(train_fn, ...)コールをtrain_fn@distributedデコレーターに置き換え、その後train_fn.distributed(...)で起動します。
  • SparkクラスターとGPUワーカーの設定を削除してください。 ノートパソコンを8xH100のアクセラレーターに接続して、代わりに gpus=8 をセットしてください。
  • デコレータで修飾された関数内にデータの読み込みを移動してください。 データの 読み込みを参照してください。
  • 既存のDDP、FSDP、またはDeepSpeedモデルのコードを保持してください。 デコレーターはこの3つすべてをサポートします。

@distributed 単一のノード上で動作するため( 制限を参照)、すべてのTorchDistributorワークロードを置き換えるわけではありません。 Spark連携に依存するワークロードはTorchDistributorに保管してください。 ローカルマシンや複数のノード間で分散トレーニングを行う場合は、 代わりにパブリックプレビューにあるAI Runtime CLIを使います。 AI ランタイム CLI を参照してください。

完全な例

以下の例は、ノートパソコンから8台のH100 GPU上でマルチレイヤーパーセプトロン(MLP)モデルを訓練しています。

  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_gpu ライブラリと distributed モジュールをインポートします。

    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 カタログ ボリュームに格納されているファイル ベースのデータの場合は、UCVolumeDatasetからのserverless_gpu.dataを使用します。このは、ローカル キャッシュを使用してファイルをストリーミングし、それらをランクとワーカー間で自動的にパーティション分割します。 分散トレーニングのチェックポイントをボリュームに保存するには、UCVolumeWriterUCVolumeReader を使用します。 AI ランタイムとモデルのチェックポイント処理に関するデータの読み込みを参照してください。

Limitations

  • 分散トレーニングはノートパソコンが接続されている単一のノード上でGPU間で実行されます。 マルチGPUの完全なトレーニングを目的とすると、8台のH100アクセラレータに接続してください。これは1ノードに8つのGPUを割り当て、 gpus=8セットします。
  • アクセラレーターの種類は一致していなければなりません。 @distributedgpu_typeを設定する場合、ノートが接続されているアクセラレータ("H100"または"A10")と一致しなければなりません。 不一致が作業負荷の失敗を引き起こします。 パラメータは任意で、省略すると自動検出されます。
  • AI RuntimeはGPU environment v4以上を推奨しています。 カスタムタイムアウト( timeout パラメータ)はGPU環境v5以上が必要です。
  • デコレーターはデフォルトで3時間後にタイムアウトします。 数秒で timeout をパスして変更するか、 timeout=None で無効にしてください。
  • 実行はノートブックのライフサイクル内で行われます。 ノートブックを終了するとランは終了します。

詳細情報