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)モデルを訓練しています。
モデルを設定し、ユーティリティ関数を定義します。
# 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_gpuライブラリとdistributedモジュールをインポートします。import serverless_gpu from serverless_gpu import distributedモデルトレーニングコードを関数でラップし、
@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()分散トレーニングはユーザー定義の引数で分散関数を呼び出して実行します。
run_train.distributed(num_epochs=3, batch_size=1)実行すると、ノートブック のセル出力に 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を使用します。このは、ローカル キャッシュを使用してファイルをストリーミングし、それらをランクとワーカー間で自動的にパーティション分割します。 分散トレーニングのチェックポイントをボリュームに保存するには、UCVolumeWriter と UCVolumeReader を使用します。 AI ランタイムとモデルのチェックポイント処理に関するデータの読み込みを参照してください。
Limitations
- 分散トレーニングはノートパソコンが接続されている単一のノード上でGPU間で実行されます。 マルチGPUの完全なトレーニングを目的とすると、8台のH100アクセラレータに接続してください。これは1ノードに8つのGPUを割り当て、
gpus=8セットします。 - アクセラレーターの種類は一致していなければなりません。
@distributedにgpu_typeを設定する場合、ノートが接続されているアクセラレータ("H100"または"A10")と一致しなければなりません。 不一致が作業負荷の失敗を引き起こします。 パラメータは任意で、省略すると自動検出されます。 - AI RuntimeはGPU environment v4以上を推奨しています。 カスタムタイムアウト(
timeoutパラメータ)はGPU環境v5以上が必要です。 - デコレーターはデフォルトで3時間後にタイムアウトします。 数秒で
timeoutをパスして変更するか、timeout=Noneで無効にしてください。 - 実行はノートブックのライフサイクル内で行われます。 ノートブックを終了するとランは終了します。
詳細情報
-
@distributedデコレーター、GPUType、Ray APIについては、Serverless GPU Python APIリファレンスドキュメントをご覧ください。 - トレーニングパイプラインをより効率的かつ回復力のあるものにするパターンについては、 パフォーマンスとレジリエンスのガイドをご覧ください。
- エンドツーエンドのトレーニングシナリオについては 、ノートブックの例を参照してください。