在筆記本中進行分散式訓練

這很重要

這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。

來自 Serverless GPU Python API 的 @distributed 裝飾器,是在 Databricks 筆記本中執行分散式訓練最方便的方式。 為你的訓練函式加上裝飾器並呼叫它,AI Runtime 就會在你的筆記本所連線節點上的所有 GPU 上執行該函式。 同樣的程式碼可以從單一 GPU 擴展到多 GPU,無需叢集配置,也無需分散式啟動器配置。

Tip

  • @distributed 裝飾器可在筆記本中,於你的節點上的每個 GPU 上執行訓練函式。
  • 它支援 PyTorch DDP、FSDP 和 DeepSpeed,並將單 GPU 程式碼以最小改動轉為多 GPU。
  • 將你的筆電連接到 8xH100 加速器,開始 gpus=8 完整的多 GPU 訓練。

Note

本頁介紹使用 Databricks 筆記本搭配無伺服器 GPU Python API 進行分散式訓練。 要從本地機器提交分散式訓練工作負載,請使用 Data Bricks 的 CLI 指令來執行 AI 執行時,這些指令目前已 在公開預覽階段。 請參閱 搭配 AI 執行階段使用 Databricks CLI。

快速入門

當你的筆電連接到無伺服器 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()

分散式訓練需要一個 8xH100 加速器,該加速器在單一節點配備 8 顆 GPU。 使用 @distributed 裝飾器時,設為 gpus=8。 這個 gpu_type 參數是可選的,會自動從你筆電連接的加速器偵測到。

每次呼叫 .distributed() 時,都會建立一個 MLflow 執行;如果已有作用中的執行,則會建立巢狀子執行,並在儲存格輸出中顯示執行連結。 完整且可執行的攻略,請參見 完整範例。

支援的架構

該 @distributed API 整合了主要的分散式訓練函式庫:

  • PyTorch 分散式資料平行(DDP):標準多 GPU 資料平行處理。
  • 全分片資料平行(FSDP):適用於大型模型的可節省記憶體訓練。
  • DeepSpeed:Microsoft 用於大型模型訓練的優化函式庫。

關於使用各函式庫的真實訓練情境,請參考 筆記本範例。

裝飾師的工作原理@distributed

當你呼叫以 .distributed() 裝飾的函式時,AI Runtime 會處理原本需要使用分散式啟動器手動設定的相關細節:

  • 序列化與扇出:函式會先序列化,然後在你要求的每個 gpus 上啟動。 每台 GPU 都會執行一個包含相同參數的函式副本。
  • 環境同步:Python 環境及其相依性在所有階級中複製,因此每個程序執行相同的程式碼。
  • Rank 環境變數:系統會針對每個程序自動設定標準變數,例如 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 上運行分散式 PyTorch,搭配 TorchDistributor ,且你的工作負載能集中在單一節點上,那麼這個 serverless_gpu@distributed API 是新深度學習工作負載的推薦替代方案。 它移除了 Spark 叢集,並提供從單一 GPU 到多 GPU 的相同程式碼路徑。

Feature serverless_gpu @distributed API 火炬分配器
基礎結構 完全無伺服器,無需叢集管理 需要一個帶有 GPU 工作者 的 Spark 叢集
設定 單一裝飾器,極簡配置 需要 Spark 叢集與 TorchDistributor 設定
架構支援 PyTorch DDP、FSDP、DeepSpeed 主要是 PyTorch DDP
數據載入 在裝飾器中,使用 Unity Catalog 磁碟區(UCVolumeDataset 用於串流檔案資料) 透過 Spark 或檔案系統

要遷移單一節點工作負載:

  • 將 train_fn.distributed(...) 上的 @distributed 呼叫替換為 TorchDistributor(...).run(train_fn, ...) 裝飾器,然後使用 train_fn 啟動。
  • 移除 Spark 叢集和 GPU 工作節點設定。 將您的筆記型電腦連接到 8xH100 加速器,然後改為設定 gpus=8。
  • 在裝飾函式中移動資料載入。 請參見資料載入。
  • 保留你現有的 DDP、FSDP 或 DeepSpeed 模型程式碼。 裝飾器支援這三者。

@distributed 在單一節點上運行(見 限制),因此不會取代所有 TorchDistributor 的工作負載。 將依賴 Spark 整合的工作負載保留在 TorchDistributor 上。 若要從本地機器或跨多個節點執行分散式訓練,請改用目前公開 預覽版的 AI 執行時 CLI。 請參閱 搭配 AI 執行階段使用 Databricks CLI。

完整範例

以下範例是在筆記型電腦中訓練一個多層感知器(MLP)模型,使用8顆H100 GPU。

  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 Catalog 磁碟區中的檔案型資料,請使用來自 UCVolumeDataset 的 serverless_gpu.data,其會以本機快取串流處理檔案,並自動在各個 rank 與 worker 之間分割檔案。 若要將分散式訓練的檢查點儲存到磁碟區,請使用 UCVolumeWriter 和 UCVolumeReader。 請參見在 AI 執行階段載入資料和使用分散式檢查點(DCP)的檢查點。

Limitations

  • 分散式訓練會在您的 notebook 所連線之單一節點上的各個 GPU 之間進行。 若要完整多 GPU 訓練,請連接 8xH100 加速器,該加速器可配置一個節點配備 8 顆 GPU,然後設定 gpus=8。
  • 加速器類型必須相符。 如果你在 @distributed 中設定 gpu_type,其值必須與你的筆記本所連線的加速器相符("H100" 或 "A10")。 不相符會導致工作負載執行失敗。 該參數為可選,省略時會自動偵測。
  • AI Runtime 推薦 GPU 環境 v4 及以上版本。 自訂逾時設定(timeout 參數)需要 v5 或以上版本的 GPU 環境。
  • 裝飾器預設會在 3 小時後逾時。 在幾秒內傳入 timeout 即可變更,或傳入 timeout=None 將其停用。
  • 執行過程在筆記本的生命週期內執行。 終止筆記本即終止該運行。

瞭解更多資訊