Adatok betöltése az AI futtatókörnyezetben

Fontos

Ez a funkció nyilvános előzetes verzióban van.

Az adat- és modelleszközök kulcsfontosságúak a mélytanulás és a képzés utáni munkaterhelések számára nagy nyelvi modellek (LLM-ek) és látás-nyelvi modellek (VLM-ek) esetében. Az AI Runtime segítségével minden adat- és modelleszközhöz a Unity Catalogon keresztül férünk hozzá:

  • Unity Catalog kötetek: elsősorban nagy adathalmazokhoz és strukturálatlan fájlokhoz, beleértve képeket, hangokat és szöveget is.
  • Unity katalógus táblák: strukturált és táblázatos adatokhoz használják, elérhető a Spark Connecten.

A köteteket és táblázatokat regisztrálni kell a Unity Catalog-ban, és elérhetővé kell tenni a felhasználó vagy szolgáltatási vezető számára.

Unity Catalog kötet strukturálatlan adatokhoz

A Unity Catalog kötetek szabályozott hozzáférést biztosítanak nem táblázatos adatokhoz bármilyen formátumban, beleértve a strukturált, félig strukturált és strukturálatlan adatokat is. Az AI Runtime-ban a kötetek az elsődleges mechanizmusok a nagy adathalmazokhoz, szöveghez, modell eszközökhöz és modellellenőrzőpontokhoz való hozzáféréshez.

A felhasználók ismerős fájlrendszer-műveletek segítségével listázhatják, olvashatják és írhatnak fájlokat Unity Catalog köteteiben, hasonlóan ahhoz, mint a helyi lemezen lévő fájlokkal való munkához:

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!")

Hasonlóképpen, a shell műveletek ugyanígy működnek:

%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

A Unity Catalog kötetek néhány jellemzője miatt jól alkalmasak gépi tanulási feladatokhoz:

  • Elosztott tárolás: A Unity Catalog elosztott tárhely támogatásával rendelkezik, lehetővé téve az AI Runtime munkaterhelések számára, hogy adatokat olvassanak és írjanak, valamint modellezzék az eszközöket a platformon keresztül, mind jegyzetfüzetekről, mind CLI-alapú terhelésekről.
  • ML hozzáférési mintákra optimalizálva: Az alapvető tárolási és hozzáférési útvonalak a gyakori ML munkaterhelésekhez vannak optimalizálva, különösen nagy fájlokhoz, amelyek sorozatos olvasási és írási folyamatokat tartalmaznak. Ez teszi a Unity Catalog-t jól alkalmassá az adatbetöltés, modelleszköz betöltésének és modellellenőrzőpontok írásának képzésére.
  • Fájlrendszer-szerű hozzáférés: A felhasználók ismerős fájlrendszer-műveletek segítségével listázhatják, olvashatják és írhatnak fájlokat Unity Catalog kötetekben, hasonlóan ahhoz, amikor helyi lemezen dolgoznak fájlokkal.

Az automatikus háttérben történő commit-ek miatt a felhasználók következetes hozzáférést várhatnak a Unity Catalog kötetadataihoz:

  • Írás: Az AI Runtime automatikusan véglegesíti az írási műveleteket, így a módosítások láthatóvá válnak a ugyanahhoz a Unity Catalog-kötethez hozzáférő más alkalmazások és munkaterhelések számára.
  • Olvasás: Az AI Runtime automatikusan észleli a kötetben bekövetkező változásokat, anélkül hogy a felhasználónak bármilyen kifejezett frissítési vagy szinkronizálási műveletet kellene végeznie.

Hanghangerő teljesítmény

Ahogy említettük, a Unity Catalog köteteket elosztott tárhely támogatja, és nagy fájlokra optimalizálva vannak egymás utáni olvasásokkal és írással.

Néhány tipp segíthet abban, hogy a legjobb teljesítményt érd el az AI Runtime-tól:

  • Összefűzd az adatokat nagyobb fájlokká: Lehetőség szerint összevond az adatokat kevesebb, nagyobb fájlba, körülbelül 1 GiB-től 10 GiB-ig fájlonként. Ez lehetővé teszi az AI Runtime számára, hogy agresszíven előre letöltse az adatokat, és automatikusan elérje a közel optimális szekvenciális olvasási teljesítményt.

  • Kis fájlos munkaterheléseknél használj helyi lemezt: Ha a munkaterhelés sok apró fájlt tartalmaz, fontold meg, hogy a fájlokat párhuzamos másolatokkal másolod helyi lemezre feldolgozás előtt. Ez csökkentheti a sok kis fájl ismétlődő elérésének megterhelését a köteten keresztül.

    # 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
    
  • Gépi tanulási terhelésekhez is használhatod UCVolumeDataset . A fent leírt optimalizálásokat tartalmazza, hogy hatékony adathozzáférést és betöltést biztosítson a Unity Catalog köteteiből. További részletekért tekintse meg az alábbi szakaszokat.

Strukturálatlan adatok betöltése ezzel: UCVolumeDataset

Strukturálatlan adatok, például a Unity Catalog-kötetekben tárolt képek, hang- és szövegfájlok esetén használja UCVolumeDataset a serverless_gpu.data csomagból. UCVolumeDataset egy PyTorch IterableDataset, amely első hozzáféréskor minden fájlt átmásol a kötetről egy gyors helyi gyorsítótárba, és a gyorsítótárazott helyi fájl elérési útját adja vissza. Kezeli azokat a teljesítmény- és terjesztési problémákat, amelyeket egyébként kézzel valósítana meg:

  • Helyi gyorsítótárazás. A rendszer az első hozzáféréskor a fájlokat a FUSE csatolásból egy helyi gyorsítótár-könyvtárba másolja, majd ezt követően a gyorsítótárból szolgálja ki őket, így a több epochon át tartó tanítás során nem olvassa újra a kötetet.
  • Automatikus particionálás. Az inicializáláskor torch.distributed a fájlok particionálva lesznek a rangok között, majd tovább vannak osztva a feldolgozók között DataLoader , így minden (rank, worker) pár egy nem egymást átfedő szeletet kap további beállítás nélkül.

Megjegyzés:

UCVolumeDataset és serverless_gpu.data.DataLoader5-ös vagy újabb GPU-környezetet igényel.

UCVolumeDataset nyers helyi fájlelérési útvonalakat eredményez. Ha ezeket a fájlokat tenzorokká szeretné dekódolni, foglalja egy második IterableDataset elembe, amely felhasználja az útvonalfolyamot, és alkalmazza a feldolgozási logikáját. Ez elkülöníti az I/O-t és az elemzési szempontokat.

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)

A wrapper már gyorsítótárazott helyi elérési utakat kap, így az elemzési lépés soha nem fér hozzá a kötethez. További burkokat is összefűzhet augmentáláshoz, tokenizáláshoz vagy szűréshez.

Az optimális teljesítmény érdekében párosítsa a(z) UCVolumeDataset elemet inkább a(z) serverless_gpu.data.DataLoader-gyel, mint a standard PyTorch DataLoader-vel. AI Runtime I/O-ra van hangolva, és egyszerre tölti be és gyorstárázza a fájlokat, miközben a GPU számításokat végez.

Ellenőrzőpont-modellek mennyiségeken

Ahhoz, hogy ellenőrizd a modelledet, hogy a legfrissebb pillanatképről folytathasd a képzést vagy vissza tudd állni egy összeomlásból, használhatod a Unity Catalog köteteket, akárcsak egy helyi fájlrendszert.

A Databricks azt javasolja, hogy elosztott ellenőrzőpontot (DCP) alkalmazzunk jobb teljesítmény érdekében mind egy-GPU-s, mind több GPU-s munkaterhelésen. Lásd a Gyors, hibaellenes PyTorch képzést az AI Runtime-ról a Databricks mérnöki blogon.

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"],
)

A monolitikus torch.save megközelítés is működik.

  • Egy-GPU modellellenőrzőpont esetén,

    # 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"])
    
  • A torchrun-on indított elosztott képzés esetén

    # 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"])
    

Táblázatos adatok betöltése

A Spark Connect használatával táblázatos gépi tanulási adatokat tölthet be Delta-táblákból.

Az egycsomópontos betanításhoz az Apache Spark DataFrame-eket a PySpark metódussaltoPandas() pandas DataFrame-ekre konvertálhatja, majd opcionálisan NumPy formátumra konvertálhatja a PySpark metódussalto_numpy().

Megjegyzés:

A Spark Connect elhalasztja az elemzést és a névfeloldást a végrehajtási időre, ami megváltoztathatja a kód viselkedését. Lásd a Spark Connect és a Klasszikus Spark összehasonlítása című témakört.

A Spark Connect a legtöbb PySpark API-t támogatja, beleértve a Spark SQL-t, a Sparkon futó Pandas API-t, a strukturált streamelést és az MLlib-t (DataFrame-alapú). Tekintse meg a PySpark API referenciadokumentációját a legújabb támogatott API-khoz.

További korlátozásokért lásd a kiszolgáló nélküli számítási korlátozásokat.

Nagy Delta táblák betöltése Unity Catalog kötetekkel

Delta nagytáblák esetén, amelyek túl nagyok ahhoz, hogy a toPandas() funkcióval átalakíthatók legyenek, exportálja az adatokat egy Unity Catalog-kötetbe, és töltse be közvetlenül a PyTorch vagy a Hugging Face használatával.

# 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")

Ez a megközelítés elkerüli a Spark terhelését a betanítás során, és jól működik az egy GPU-s és az elosztott betanítási munkafolyamatok esetében is.