Ładowanie danych w środowisku uruchomieniowym sztucznej inteligencji

Ważna

Ta funkcja jest dostępna w publicznej wersji testowej.

Dane i zasoby modelowe mają kluczowe znaczenie w uczeniu głębokim oraz w obciążeniach posttreningowych związanych z dużymi modelami językowymi (LLM) i modelami wizyjno-językowymi (VLM). Dzięki AI Runtime wszystkie dane i zasoby modelu są dostępne przez Unity Catalog:

  • Tomy katalogu Unity: używane głównie do dużych zbiorów danych i plików nieustrukturyzowanych, w tym obrazów, dźwięku i tekstu.
  • Tabele katalogu Unity: używane do danych strukturalnych i tabelarowych, dostępne przez Spark Connect.

Twoje tomy i tabele muszą być zarejestrowane w Unity Catalog i dostępne dla użytkownika lub podmiotu usługi.

Wolumin katalogu Unity dla danych nieustrukturyzowanych

Tomy Unity Catalog zapewniają kontrolowany dostęp do danych nietabelarnych w dowolnym formacie, w tym danych strukturalnych, półstrukturalnych i nieustrukturyzowanych. W AI Runtime wolumeny są głównym mechanizmem dostępu do dużych zbiorów danych, tekstów, zasobów modelu i punktów kontrolnych modelu.

Użytkownicy mogą wymieniać, odczytywać i zapisywać pliki w woluminach Unity Catalog, korzystając z znanych operacji systemu plików, podobnie jak pracując z plikami na lokalnym dysku:

import os

dir_path = "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir"
file_path = os.path.join(dir_path, "test_file")

os.makedirs(dir_path, exist_ok=True)

# Write to the file
with open(file_path, "w") as file:
    file.write("Hello, World!")

Podobnie operacje powłoki działają w ten sam sposób:

%sh ls -l /Volumes/<catalog-name>/<schema-name>/<volume-name>
%sh mkdir -p /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir
%sh touch /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/test_file

Kilka cech tomów Unity Catalog sprawia, że są one dobrze przystosowane do zadań związanych z uczeniem maszynowym:

  • Rozproszone przechowywanie: Unity Catalog jest wspierany przez rozproszone przechowywanie, co pozwala obciążeniom AI Runtime odczytywać i zapisywać dane oraz modelować zasoby na całej platformie, zarówno z notebooków, jak i obciążeń opartych na CLI.
  • Optymalizacja pod wzorce dostępu ML: Podstawowe ścieżki przechowywania i dostępu są optymalizowane pod typowe obciążenia ML, szczególnie duże pliki z sekwencyjnymi odczytami i zapisami. Dzięki temu Unity Catalog doskonale nadaje się do treningowego ładowania danych, ładowania zasobów modeli oraz zapisywania punktów kontrolnych modelu.
  • Dostęp podobny do systemu plików: Użytkownicy mogą wymieniać, czytać i zapisywać pliki w woluminach Unity Catalog, korzystając z znanych operacji systemu plików, podobnie jak praca z plikami na lokalnym dysku.

Dzięki automatycznym zapisom w tle użytkownicy mogą liczyć na stały dostęp do danych woluminów Unity Catalog:

  • Zapisy: AI Runtime automatycznie zatwierdza zapisy, czyniąc zmiany widocznymi dla innych aplikacji i obciążeń korzystających z tego samego woluminu katalogu Unity.
  • Odczytuje: AI Runtime automatycznie odbiera zmiany w głośności bez konieczności stosowania wyraźnego odświeżania czy synchronizacji przez użytkownika.

Dostosuj działanie głośności

Jak wspomniano, tomy katalogu Unity są wspierane przez rozproszoną pamięć masową i zoptymalizowane dla dużych plików z sekwencyjnymi odczytami i zapisami.

Kilka wskazówek pomoże Ci uzyskać najlepszą wydajność z AI Runtime:

  • Łącz dane do większych plików: Gdy to możliwe, konsoliduj dane w mniej, ale większe pliki, około 1 do 10 GiB na plik. Pozwala to AI Runtime na agresywne pobieranie danych i automatyczne osiągnięcie niemal optymalnej wydajności odczytu sekwencyjnego.

  • Do zadań z małymi plikami używaj lokalnego dysku: Jeśli Twoje obciążenie obejmuje wiele małych plików, rozważ kopiowanie plików na lokalny dysk za pomocą kopii równoległych przed ich przetwarzaniem. Może to zmniejszyć narzut związany z wielokrotnym uzyskiwaniem dostępu do wielu małych plików przez wolumín.

    # Recommended using parallel copy (256 concurrency in this example, you can tune)
    #
    # This takes only 22 seconds to copy 15,375 150KiB small image files.
    %sh cd /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ && find . -type f -print0 | xargs -0 -P 256 -I {} cp --parents "{}" /tmp/
    
    
    # !!! Avoid doing this !!!
    #
    # Because the files are copied in serial, this copies the same 15,375 150KiB small image files much more slowly.
    # %sh cp -r /Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/* /tmp
    
  • Możesz użyć tego UCVolumeDataset do swoich zadań związanych z uczeniem maszynowym. Zawiera opisane powyżej optymalizacje, aby zapewnić efektywny dostęp do danych i ładowanie z woluminów katalogu Unity. Aby uzyskać więcej informacji, zobacz poniższe sekcje.

Załaduj dane niestrukturalne za pomocą UCVolumeDataset

W przypadku danych nieustrukturyzowanych, takich jak obrazy, pliki audio i pliki tekstowe przechowywane w woluminach Unity Catalog, użyj UCVolumeDataset z pakietu serverless_gpu.data. UCVolumeDataset to element PyTorch IterableDataset, który przy pierwszym dostępie kopiuje każdy plik z woluminu do szybkiej lokalnej pamięci podręcznej i zwraca ścieżkę do zbuforowanego lokalnego pliku. Obsługuje ona problemy z wydajnością i dystrybucją, które w przeciwnym razie należy zaimplementować ręcznie:

  • Buforowanie lokalne. Pliki są kopiowane z punktu montowania FUSE do lokalnego katalogu pamięci podręcznej przy pierwszym dostępie, a następnie udostępniane z pamięci podręcznej, więc trening wieloepokowy nie powoduje ponownego odczytu woluminu.
  • Automatyczne partycjonowanie. Po zainicjowaniu torch.distributed pliki są dzielone między rangi, a następnie między procesy robocze DataLoader, dzięki czemu każda para (rank, worker) otrzymuje niepokrywający się fragment bez dodatkowej konfiguracji.

Uwaga / Notatka

UCVolumeDataset i serverless_gpu.data.DataLoader wymagają środowiska procesora GPU 5 lub nowszego.

UCVolumeDataset zwraca nieprzetworzone ścieżki plików lokalnych. Aby zdekodować te pliki do tensorów, opakuj je w drugi element IterableDataset, który korzysta ze strumienia ścieżek i stosuje logikę parsowania. Dzięki temu kwestie związane z I/O i parsowaniem pozostają oddzielone.

from serverless_gpu.data import UCVolumeDataset
from torch.utils.data import IterableDataset
from PIL import Image
import torchvision.transforms.functional as TF

class ImageDataset(IterableDataset):
    """Decodes each cached file path from UCVolumeDataset into a tensor."""

    def __init__(self, path_dataset: UCVolumeDataset):
        self._path_dataset = path_dataset

    def __iter__(self):
        for local_path in self._path_dataset:
            image = Image.open(local_path).convert("RGB")
            yield TF.to_tensor(image)

path_dataset = UCVolumeDataset("/Volumes/catalog/schema/my_volume/images")
dataset = ImageDataset(path_dataset)

Warstwa opakowująca otrzymuje lokalne ścieżki zapisane już w pamięci podręcznej, więc etap parsowania nigdy nie uzyskuje dostępu do woluminu. Możesz łączyć w łańcuch dodatkowe opakowania do rozszerzania, tokenizacji lub filtrowania.

Aby uzyskać optymalną wydajność, używaj UCVolumeDataset w połączeniu z serverless_gpu.data.DataLoader, zamiast domyślnego DataLoader z PyTorch. Jest dostrojony pod AI Runtime I/O i pobiera oraz zapisuje pliki jednocześnie podczas obliczeń GPU.

Modele punktów kontrolnych na woluminach

Aby zapisywać punkty kontrolne modelu, dzięki czemu można wznowić trenowanie od najnowszej migawki lub odzyskać stan po awarii, możesz używać woluminów Unity Catalog tak jak lokalnego systemu plików.

Databricks zaleca stosowanie rozproszonego punktu kontrolnego (DCP) dla lepszej wydajności zarówno na obciążeniach z pojedynczym GPU, jak i z wieloma GPU. Zobacz szybkie, odporne na błędy szkolenia PyTorch na AI Runtime na blogu inżynierskim Databricks.

import torch.distributed.checkpoint as dcp
from torch.distributed.checkpoint.state_dict import get_state_dict, set_state_dict

import serverless_gpu

checkpoint_path = "/Volumes/my-catalog/my-schema/my-volume/checkpoints/step_1000"

# Save
model_sd, optim_sd = get_state_dict(model, optimizer)
state_dict = {"model": model_sd, "optim": optim_sd, "step": 1000}
dcp.async_save(
    state_dict,
    storage_writer=serverless_gpu.data.UCVolumeWriter(checkpoint_path))

# Load
model_sd, optim_sd = get_state_dict(model, optimizer)
state_dict = {"model": model_sd, "optim": optim_sd}
dcp.load(
    state_dict,
    storage_reader=serverless_gpu.data.UCVolumeReader(checkpoint_path))

set_state_dict(
    model,
    optimizer,
    model_state_dict=state_dict["model"],
    optim_state_dict=state_dict["optim"],
)

Monolityczne podejście torch.save też działa.

  • Dla checkpointingu modelu z jednym GPU,

    # The monolithic torch.save approach for single GPU chip
    
    # Save
    torch.save({"model": model.state_dict(), "opt": optimizer.state_dict()},
               "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt")
    
    
    # Load
    ckpt = torch.load(
        "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt",
        weights_only=True)
    model.load_state_dict(ckpt["model"])
    optimizer.load_state_dict(ckpt["opt"])
    
  • W przypadku treningu rozproszonego uruchamianego za pomocą torchrun,

    # The monolithic torch.save approach for multi-GPU distributed training.
    # This snippet assumes your launcher has already called
    # dist.init_process_group(...).
    
    import os
    import torch.distributed as dist
    
    # Save only on rank 0.
    if dist.get_rank() == 0:
        torch.save({"model": model.state_dict(), "opt": optimizer.state_dict()},
                   "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt")
    
    # Wait for rank 0 to finish writing before any rank reads.
    dist.barrier()
    
    # Load on ALL ranks (map to current rank's local GPU).
    local_rank = int(os.environ["LOCAL_RANK"])
    ckpt = torch.load(
        "/Volumes/<catalog-name>/<schema-name>/<volume-name>/sub-dir/ckpt.pt",
        map_location=f"cuda:{local_rank}",
        weights_only=True)
    model.load_state_dict(ckpt["model"])
    optimizer.load_state_dict(ckpt["opt"])
    

Ładowanie danych tabelarycznych

Użyj programu Spark Connect, aby załadować tabelaryczne dane uczenia maszynowego z tabel delta.

W przypadku trenowania z jednym węzłem można przekonwertować ramki danych platformy Apache Spark na ramki danych biblioteki pandas przy użyciu metody toPandas(), a następnie opcjonalnie przekonwertować na format NumPy przy użyciu metody to_numpy().

Uwaga / Notatka

Spark Connect odracza analizę i rozpoznawanie nazw do czasu kompilacji, co może wpłynąć na sposób działania twojego kodu. Zobacz Porównanie programu Spark Connect z modelem klasycznym platformy Spark.

Narzędzie Spark Connect obsługuje większość interfejsów API PySpark, w tym Spark SQL, interfejs API Pandas na platformie Spark, uporządkowane przesyłanie strumieniowe oraz bibliotekę MLlib (opartą na ramkach danych). Aby uzyskać najnowsze obsługiwane interfejsy API, zobacz dokumentację interfejsu API PySpark .

Aby uzyskać informacje o innych ograniczeniach, zobacz Ograniczenia obliczeniowe bezserwerowe.

Ładuj duże tabele Delta za pomocą wolumenów Unity Catalog

W przypadku dużych tabel Delta, które są zbyt duże do konwersji za pomocą toPandas(), wyeksportuj dane do woluminu Unity Catalog i załaduj je bezpośrednio z użyciem PyTorch lub Hugging Face.

# Step 1: Export the Delta table to Parquet files in a UC volume
output_path = "/Volumes/catalog/schema/my_volume/training_data"
spark.table("catalog.schema.my_table").write.mode("overwrite").parquet(output_path)
# Step 2: Load the exported data directly using Hugging Face datasets
from datasets import load_dataset

dataset = load_dataset("parquet", data_files="/Volumes/catalog/schema/my_volume/training_data/*.parquet")

Takie podejście pozwala uniknąć narzutów platformy Spark podczas trenowania i działa dobrze zarówno w przypadku procesów trenowania pojedynczego GPU, jak i rozproszonej pracy trenowania.